Table of contents
Open Table of contents
Consumer Callbacks
After some emails are sent successfully, follow-up logic needs to run, such as updating a database. At this point, Message Server becomes a producer that sends callback messages to Kafka, while Business Server becomes a consumer of those messages.
How Should Callback Messages Be Abstracted?
Callback logic depends on the business scenario. How can we support different callback logic while keeping a consistent message format? We use reflection.
@JsonDeserialize(builder = CallbackMetaData.CallbackMetaDataBuilder.class)
@JsonInclude(JsonInclude.Include.NON_NULL)
@Getter
@Builder
@ToString
public class CallbackMetaData implements Serializable {
@JsonProperty("id")
private String messageId;
/**
* This should be hostName: InetAddress.getLocalHost().getHostName().
*/
@JsonProperty("serverId")
private String serverId;
@JsonProperty("className")
private String className;
/**
* this string is the json string of the instance of the class, generated by Jackson.
* for example:
* className instance = new className();
* objectMapper.writeValueAsString(instance);
*/
@JsonProperty("instanceJsonStr")
private String instanceJsonStr;
@JsonProperty("methodName")
private String methodName;
@JsonProperty("arguments")
private Object[] arguments;
}
This metadata is sent to Message Server as part of an email message. After sending the email, Message Server checks whether callback information is present. If so, it sends CallbackMetaData to the relevant topic.
Why Set serverId?
There are two reasons:
- The code is deployed through rolling updates. For compatibility, a callback must return to the same Business Server that included it in the original message.
- Initially, we planned to use one topic subscribed to by all Business Servers, with different consumer groups to broadcast messages. Each server would compare the message’s serverId with its own: consume it if they matched, or commit the offset directly otherwise. We later realized this made every server consume messages continuously, including messages not meant for it. The improved design gives every server its own callback topic and consumes only that topic. The topic name is callback-serverId.
serverId uses the machine’s hostName rather than its IP address, because the IP may change.
How Are Callback Messages Consumed?
I will not repeat the other consumer code; see the previous article for details. The consumption logic itself is only three lines:
ConsumerRecords<String, CallbackMetaData> records = consumer.poll(Duration.ofSeconds(10));
records.forEach(each -> {
Class<?> destClass;
try {
// Core consumption logic: invoke the target method through reflection.
destClass = Class.forName(each.value().getClassName());
Object instance = objectMapper.readValue(each.value().getInstanceJsonStr(), destClass);
MethodUtils.invokeMethod(instance, true, each.value().getMethodName(), each.value().getArguments());
} catch (Exception e) {
e.printStackTrace();
}
});
Scheduled Retries
The previous two articles explained that producer or consumer messages still failing after N retries are stored in a database for later retries through a scheduled job. To reduce the burden on business servers, Message Server handles all failed-message retries.
Failed-Message Table Design
@Getter
@Setter
@ToString
@EqualsAndHashCode(of = {"messageId", "failedPhrases"})
public class MessageFailedEntity {
/**
* Primary key
*/
private Long id;
/**
* Message ID
*/
private String messageId;
/**
* Message contents in JSON format
*/
private String messageContentJsonFormat;
/**
* Message type
* EMAIL indicates an email message.
* EMAIL_CALLBACK indicates an email callback.
*
*/
private MessageType messageType;
/**
* Stage at which the message failed:
* PRODUCER indicates failure while producing the message.
* CONSUMER indicates failure while consuming the message.
*/
private MessageFailedPhase messageFailedPhase;
/**
* Exception stack trace from the failure
*/
private String failedReason;
/**
* Message retry count
*/
private Integer retryCount;
/**
* Message retry status
* 0 indicates retry failure.
* 1 indicates retry success.
*/
private Integer retryStatus;
/**
* Timestamp
*/
private LocalDateTime lastUpdateTime;
}
Retry Logic Design
The retry approach is simple:
- Query a batch of records from the failed-message table, perhaps 100 or 10000 at a time, depending on the scenario.
- Deserialize the JSON into the corresponding object based on the message type, then use the appropriate producer to send it to Kafka.
- If the failure stage is PRODUCER, update the record to mark the retry successful after a successful resend.
- If the failure stage is CONSUMER, a successful resend only updates the retry count. The relevant consumer updates whether the retry succeeded, because a CONSUMER retry succeeds only once consumption succeeds.
- Set a maximum retry count and stop retrying when it is exceeded.
- With multiple Message Servers, a distributed lock can ensure that only one server executes the scheduled retry job at a time. The main purpose is to prevent concurrent jobs from retrieving the same batch from the database. A table flag indicating whether a record is being retried can achieve the same purpose; choose based on your scenario.
Together with the previous two articles, this handles the exception cases that may occur throughout the messaging system.
Understanding Rebalancing
Kafka: The Definitive Guide > Chapter 4, Section 1.
Moving partition ownership from one consumer to another is called a rebalance.
When I first encountered rebalancing, I wondered: if a consumer is still processing a message when Kafka needs to rebalance, what happens to its business logic? Could it be interrupted midway through consumption? If so, that would make idempotency considerably harder to achieve.
With these questions in mind, I searched for information and found a Confluent blog post explaining the rebalance process in detail: article.
The following material comes from that article.
Suppose we have an existing consumer group with a set assignment of topic-partitions to consumers. This consumer group consists of a number of consumers, each with a member id as well as a group leader (usually the consumer that was first to join the group). A new consumer comes along and requests to join the consumer group by sending a request of JoinGroup to the Group Coordinator along with the topics it would like to subscribe to.
The Group Coordinator kicks off the rebalance by telling all current members to issue their own JoinGroup requests. This is done as part of the response to the heartbeat that consumers send to the Group Coordinator to tell it they’re still alive and well.
Each consumer in the group has max.poll.interval.ms to wrap up their current processing and send their JoinGroup request, at which point the world is stopped. With all of the JoinGroup requests, the Group Coordinator knows all of the consumers in the group and which topics should be part of the consumer group. It sends JoinResponses to the members, chooses a leader from amongst the members, and leaves the leader to compute the partition assignments.
All group members respond with a SyncGroup request. The group leader sends its partition assignments along with its request.
At this point, the Group Coordinator can send its SyncResponse to each consumer confirming their assigned topic-partitions.
Finally, consumers acknowledge their assignments and processing can resume. The world is no longer stopped
The second highlighted section answered my question: before rebalancing, it waits for each consumer to finish its processing logic.
Understanding Rebalancing through Logs
The following logs are from a local rebalance. Compare them with the steps above to deepen your understanding.
2023-11-08T02:23:04.180-0500 kafka-coordinator-heartbeat-thread | low-priority-email-group org.apache.kafka.clients.consumer.internals.ConsumerCoordinator WARN: [Consumer instanceId=2, clientId=consumer-low-priority-email-group-2, groupId=low-priority-email-group] consumer poll timeout has expired. This means the time between subsequent calls to poll() was longer than the configured max.poll.interval.ms, which typically implies that the poll loop is spending too much time processing messages. You can address this either by increasing max.poll.interval.ms or by reducing the maximum size of batches returned in poll() with max.poll.records.
2023-11-08T02:23:04.180-0500 kafka-coordinator-heartbeat-thread | low-priority-email-group org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=2, clientId=consumer-low-priority-email-group-2, groupId=low-priority-email-group] Resetting generation and member id due to: consumer pro-actively leaving the group
2023-11-08T02:23:04.180-0500 kafka-coordinator-heartbeat-thread | low-priority-email-group org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=2, clientId=consumer-low-priority-email-group-2, groupId=low-priority-email-group] Request joining group due to: consumer pro-actively leaving the group
2023-11-08T02:32:21.607-0500 consumer-low-priority-email-pool-3 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=4, clientId=consumer-low-priority-email-group-4, groupId=low-priority-email-group] Request joining group due to: group is already rebalancing
2023-11-08T02:32:21.620-0500 consumer-low-priority-email-pool-3 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=4, clientId=consumer-low-priority-email-group-4, groupId=low-priority-email-group] Revoke previously assigned partitions low-priority-email-4
2023-11-08T02:32:21.621-0500 consumer-low-priority-email-pool-3 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=4, clientId=consumer-low-priority-email-group-4, groupId=low-priority-email-group] (Re-)joining group
2023-11-08T02:32:21.712-0500 consumer-low-priority-email-pool-8 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=9, clientId=consumer-low-priority-email-group-9, groupId=low-priority-email-group] Request joining group due to: group is already rebalancing
2023-11-08T02:32:21.712-0500 consumer-low-priority-email-pool-8 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=9, clientId=consumer-low-priority-email-group-9, groupId=low-priority-email-group] Revoke previously assigned partitions low-priority-email-9
2023-11-08T02:32:21.712-0500 consumer-low-priority-email-pool-8 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=9, clientId=consumer-low-priority-email-group-9, groupId=low-priority-email-group] (Re-)joining group
2023-11-08T02:32:21.723-0500 consumer-low-priority-email-pool-7 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=8, clientId=consumer-low-priority-email-group-8, groupId=low-priority-email-group] Request joining group due to: group is already rebalancing
2023-11-08T02:32:21.723-0500 consumer-low-priority-email-pool-7 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=8, clientId=consumer-low-priority-email-group-8, groupId=low-priority-email-group] Revoke previously assigned partitions low-priority-email-8
2023-11-08T02:32:21.723-0500 consumer-low-priority-email-pool-7 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=8, clientId=consumer-low-priority-email-group-8, groupId=low-priority-email-group] (Re-)joining group
2023-11-08T02:32:21.739-0500 consumer-low-priority-email-pool-9 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=10, clientId=consumer-low-priority-email-group-10, groupId=low-priority-email-group] Request joining group due to: group is already rebalancing
2023-11-08T02:32:21.739-0500 consumer-low-priority-email-pool-9 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=10, clientId=consumer-low-priority-email-group-10, groupId=low-priority-email-group] Revoke previously assigned partitions low-priority-email-1
2023-11-08T02:32:21.739-0500 consumer-low-priority-email-pool-9 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=10, clientId=consumer-low-priority-email-group-10, groupId=low-priority-email-group] (Re-)joining group
2023-11-08T02:32:21.741-0500 consumer-low-priority-email-pool-4 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=5, clientId=consumer-low-priority-email-group-5, groupId=low-priority-email-group] Request joining group due to: group is already rebalancing
2023-11-08T02:32:21.741-0500 consumer-low-priority-email-pool-4 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=5, clientId=consumer-low-priority-email-group-5, groupId=low-priority-email-group] Revoke previously assigned partitions low-priority-email-5
2023-11-08T02:32:21.741-0500 consumer-low-priority-email-pool-4 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=5, clientId=consumer-low-priority-email-group-5, groupId=low-priority-email-group] (Re-)joining group
2023-11-08T02:32:21.745-0500 consumer-low-priority-email-pool-5 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=6, clientId=consumer-low-priority-email-group-6, groupId=low-priority-email-group] Request joining group due to: group is already rebalancing
2023-11-08T02:32:21.745-0500 consumer-low-priority-email-pool-5 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=6, clientId=consumer-low-priority-email-group-6, groupId=low-priority-email-group] Revoke previously assigned partitions low-priority-email-6
2023-11-08T02:32:21.745-0500 consumer-low-priority-email-pool-5 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=6, clientId=consumer-low-priority-email-group-6, groupId=low-priority-email-group] (Re-)joining group
2023-11-08T02:32:21.752-0500 consumer-low-priority-email-pool-2 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=3, clientId=consumer-low-priority-email-group-3, groupId=low-priority-email-group] Request joining group due to: group is already rebalancing
2023-11-08T02:32:21.753-0500 consumer-low-priority-email-pool-2 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=3, clientId=consumer-low-priority-email-group-3, groupId=low-priority-email-group] Revoke previously assigned partitions low-priority-email-3
2023-11-08T02:32:21.753-0500 consumer-low-priority-email-pool-2 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=3, clientId=consumer-low-priority-email-group-3, groupId=low-priority-email-group] (Re-)joining group
2023-11-08T02:32:21.903-0500 consumer-low-priority-email-pool-0 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=1, clientId=consumer-low-priority-email-group-1, groupId=low-priority-email-group] Request joining group due to: group is already rebalancing
2023-11-08T02:32:21.903-0500 consumer-low-priority-email-pool-0 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=1, clientId=consumer-low-priority-email-group-1, groupId=low-priority-email-group] Revoke previously assigned partitions low-priority-email-0
2023-11-08T02:32:21.903-0500 consumer-low-priority-email-pool-0 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=1, clientId=consumer-low-priority-email-group-1, groupId=low-priority-email-group] (Re-)joining group
2023-11-08T02:32:21.919-0500 consumer-low-priority-email-pool-6 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=7, clientId=consumer-low-priority-email-group-7, groupId=low-priority-email-group] Request joining group due to: group is already rebalancing
2023-11-08T02:32:21.919-0500 consumer-low-priority-email-pool-6 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=7, clientId=consumer-low-priority-email-group-7, groupId=low-priority-email-group] Revoke previously assigned partitions low-priority-email-7
2023-11-08T02:32:21.919-0500 consumer-low-priority-email-pool-6 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=7, clientId=consumer-low-priority-email-group-7, groupId=low-priority-email-group] (Re-)joining group
2023-11-08T02:32:21.920-0500 consumer-low-priority-email-pool-0 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=1, clientId=consumer-low-priority-email-group-1, groupId=low-priority-email-group] Successfully joined group with generation Generation{generationId=12, memberId='1-f1f4f621-3bde-4e02-9621-a706320300ae', protocol='range'}
2023-11-08T02:32:21.920-0500 consumer-low-priority-email-pool-9 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=10, clientId=consumer-low-priority-email-group-10, groupId=low-priority-email-group] Successfully joined group with generation Generation{generationId=12, memberId='10-cd494fa6-1bf4-4dd8-ae28-a494b21ad247', protocol='range'}
2023-11-08T02:32:21.920-0500 consumer-low-priority-email-pool-6 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=7, clientId=consumer-low-priority-email-group-7, groupId=low-priority-email-group] Successfully joined group with generation Generation{generationId=12, memberId='7-e77860ca-7059-4a38-bfc9-7db3cc862a38', protocol='range'}
2023-11-08T02:32:21.921-0500 consumer-low-priority-email-pool-2 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=3, clientId=consumer-low-priority-email-group-3, groupId=low-priority-email-group] Successfully joined group with generation Generation{generationId=12, memberId='3-ec05c279-8059-4f03-b804-3da299f93b88', protocol='range'}
2023-11-08T02:32:21.921-0500 consumer-low-priority-email-pool-4 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=5, clientId=consumer-low-priority-email-group-5, groupId=low-priority-email-group] Successfully joined group with generation Generation{generationId=12, memberId='5-d08406a6-2bd1-4bed-aed9-5f65d7f75260', protocol='range'}
2023-11-08T02:32:21.921-0500 consumer-low-priority-email-pool-8 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=9, clientId=consumer-low-priority-email-group-9, groupId=low-priority-email-group] Successfully joined group with generation Generation{generationId=12, memberId='9-bab69ccd-a0b1-4f94-99d5-869fce905f3e', protocol='range'}
2023-11-08T02:32:21.921-0500 consumer-low-priority-email-pool-5 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=6, clientId=consumer-low-priority-email-group-6, groupId=low-priority-email-group] Successfully joined group with generation Generation{generationId=12, memberId='6-67b38de0-90e0-4f53-aaae-49c10a91a463', protocol='range'}
2023-11-08T02:32:21.921-0500 consumer-low-priority-email-pool-3 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=4, clientId=consumer-low-priority-email-group-4, groupId=low-priority-email-group] Successfully joined group with generation Generation{generationId=12, memberId='4-b24fc609-920a-4cc3-b507-2a4a1f1a568b', protocol='range'}
2023-11-08T02:32:21.921-0500 consumer-low-priority-email-pool-7 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=8, clientId=consumer-low-priority-email-group-8, groupId=low-priority-email-group] Successfully joined group with generation Generation{generationId=12, memberId='8-70cd2b5d-e17f-4346-bfcb-b170e766db39', protocol='range'}
2023-11-08T02:32:21.921-0500 consumer-low-priority-email-pool-3 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=4, clientId=consumer-low-priority-email-group-4, groupId=low-priority-email-group] Finished assignment for group at generation 12: {10-cd494fa6-1bf4-4dd8-ae28-a494b21ad247=Assignment(partitions=[low-priority-email-2]), 7-e77860ca-7059-4a38-bfc9-7db3cc862a38=Assignment(partitions=[low-priority-email-7]), 3-ec05c279-8059-4f03-b804-3da299f93b88=Assignment(partitions=[low-priority-email-3]), 6-67b38de0-90e0-4f53-aaae-49c10a91a463=Assignment(partitions=[low-priority-email-6]), 8-70cd2b5d-e17f-4346-bfcb-b170e766db39=Assignment(partitions=[low-priority-email-8]), 1-f1f4f621-3bde-4e02-9621-a706320300ae=Assignment(partitions=[low-priority-email-0, low-priority-email-1]), 4-b24fc609-920a-4cc3-b507-2a4a1f1a568b=Assignment(partitions=[low-priority-email-4]), 5-d08406a6-2bd1-4bed-aed9-5f65d7f75260=Assignment(partitions=[low-priority-email-5]), 9-bab69ccd-a0b1-4f94-99d5-869fce905f3e=Assignment(partitions=[low-priority-email-9])}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-9 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=10, clientId=consumer-low-priority-email-group-10, groupId=low-priority-email-group] Successfully synced group in generation Generation{generationId=12, memberId='10-cd494fa6-1bf4-4dd8-ae28-a494b21ad247', protocol='range'}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-0 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=1, clientId=consumer-low-priority-email-group-1, groupId=low-priority-email-group] Successfully synced group in generation Generation{generationId=12, memberId='1-f1f4f621-3bde-4e02-9621-a706320300ae', protocol='range'}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-0 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=1, clientId=consumer-low-priority-email-group-1, groupId=low-priority-email-group] Notifying assignor about the new Assignment(partitions=[low-priority-email-0, low-priority-email-1])
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-9 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=10, clientId=consumer-low-priority-email-group-10, groupId=low-priority-email-group] Notifying assignor about the new Assignment(partitions=[low-priority-email-2])
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-9 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=10, clientId=consumer-low-priority-email-group-10, groupId=low-priority-email-group] Adding newly assigned partitions: low-priority-email-2
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-0 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=1, clientId=consumer-low-priority-email-group-1, groupId=low-priority-email-group] Adding newly assigned partitions: low-priority-email-0, low-priority-email-1
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-2 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=3, clientId=consumer-low-priority-email-group-3, groupId=low-priority-email-group] Successfully synced group in generation Generation{generationId=12, memberId='3-ec05c279-8059-4f03-b804-3da299f93b88', protocol='range'}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-4 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=5, clientId=consumer-low-priority-email-group-5, groupId=low-priority-email-group] Successfully synced group in generation Generation{generationId=12, memberId='5-d08406a6-2bd1-4bed-aed9-5f65d7f75260', protocol='range'}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-2 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=3, clientId=consumer-low-priority-email-group-3, groupId=low-priority-email-group] Notifying assignor about the new Assignment(partitions=[low-priority-email-3])
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-2 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=3, clientId=consumer-low-priority-email-group-3, groupId=low-priority-email-group] Adding newly assigned partitions: low-priority-email-3
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-6 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=7, clientId=consumer-low-priority-email-group-7, groupId=low-priority-email-group] Successfully synced group in generation Generation{generationId=12, memberId='7-e77860ca-7059-4a38-bfc9-7db3cc862a38', protocol='range'}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-4 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=5, clientId=consumer-low-priority-email-group-5, groupId=low-priority-email-group] Notifying assignor about the new Assignment(partitions=[low-priority-email-5])
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-7 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=8, clientId=consumer-low-priority-email-group-8, groupId=low-priority-email-group] Successfully synced group in generation Generation{generationId=12, memberId='8-70cd2b5d-e17f-4346-bfcb-b170e766db39', protocol='range'}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-6 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=7, clientId=consumer-low-priority-email-group-7, groupId=low-priority-email-group] Notifying assignor about the new Assignment(partitions=[low-priority-email-7])
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-3 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=4, clientId=consumer-low-priority-email-group-4, groupId=low-priority-email-group] Successfully synced group in generation Generation{generationId=12, memberId='4-b24fc609-920a-4cc3-b507-2a4a1f1a568b', protocol='range'}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-6 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=7, clientId=consumer-low-priority-email-group-7, groupId=low-priority-email-group] Adding newly assigned partitions: low-priority-email-7
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-4 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=5, clientId=consumer-low-priority-email-group-5, groupId=low-priority-email-group] Adding newly assigned partitions: low-priority-email-5
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-7 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=8, clientId=consumer-low-priority-email-group-8, groupId=low-priority-email-group] Notifying assignor about the new Assignment(partitions=[low-priority-email-8])
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-8 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=9, clientId=consumer-low-priority-email-group-9, groupId=low-priority-email-group] Successfully synced group in generation Generation{generationId=12, memberId='9-bab69ccd-a0b1-4f94-99d5-869fce905f3e', protocol='range'}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-5 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=6, clientId=consumer-low-priority-email-group-6, groupId=low-priority-email-group] Successfully synced group in generation Generation{generationId=12, memberId='6-67b38de0-90e0-4f53-aaae-49c10a91a463', protocol='range'}
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-3 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=4, clientId=consumer-low-priority-email-group-4, groupId=low-priority-email-group] Notifying assignor about the new Assignment(partitions=[low-priority-email-4])
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-7 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=8, clientId=consumer-low-priority-email-group-8, groupId=low-priority-email-group] Adding newly assigned partitions: low-priority-email-8
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-8 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=9, clientId=consumer-low-priority-email-group-9, groupId=low-priority-email-group] Notifying assignor about the new Assignment(partitions=[low-priority-email-9])
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-3 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=4, clientId=consumer-low-priority-email-group-4, groupId=low-priority-email-group] Adding newly assigned partitions: low-priority-email-4
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-5 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=6, clientId=consumer-low-priority-email-group-6, groupId=low-priority-email-group] Notifying assignor about the new Assignment(partitions=[low-priority-email-6])
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-8 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=9, clientId=consumer-low-priority-email-group-9, groupId=low-priority-email-group] Adding newly assigned partitions: low-priority-email-9
2023-11-08T02:32:21.922-0500 consumer-low-priority-email-pool-5 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=6, clientId=consumer-low-priority-email-group-6, groupId=low-priority-email-group] Adding newly assigned partitions: low-priority-email-6
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-9 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=10, clientId=consumer-low-priority-email-group-10, groupId=low-priority-email-group] Setting offset for partition low-priority-email-2 to the committed offset FetchPosition{offset=13436, offsetEpoch=Optional[8], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=8}}
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-0 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=1, clientId=consumer-low-priority-email-group-1, groupId=low-priority-email-group] Setting offset for partition low-priority-email-1 to the committed offset FetchPosition{offset=13301, offsetEpoch=Optional[8], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=8}}
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-0 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=1, clientId=consumer-low-priority-email-group-1, groupId=low-priority-email-group] Setting offset for partition low-priority-email-0 to the committed offset FetchPosition{offset=24087, offsetEpoch=Optional[8], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=8}}
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-2 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=3, clientId=consumer-low-priority-email-group-3, groupId=low-priority-email-group] Setting offset for partition low-priority-email-3 to the committed offset FetchPosition{offset=13299, offsetEpoch=Optional[2], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=2}}
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-3 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=4, clientId=consumer-low-priority-email-group-4, groupId=low-priority-email-group] Setting offset for partition low-priority-email-4 to the committed offset FetchPosition{offset=13352, offsetEpoch=Optional[2], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=2}}
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-7 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=8, clientId=consumer-low-priority-email-group-8, groupId=low-priority-email-group] Setting offset for partition low-priority-email-8 to the committed offset FetchPosition{offset=13190, offsetEpoch=Optional[2], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=2}}
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-8 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=9, clientId=consumer-low-priority-email-group-9, groupId=low-priority-email-group] Setting offset for partition low-priority-email-9 to the committed offset FetchPosition{offset=13151, offsetEpoch=Optional[2], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=2}}
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-5 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=6, clientId=consumer-low-priority-email-group-6, groupId=low-priority-email-group] Setting offset for partition low-priority-email-6 to the committed offset FetchPosition{offset=13211, offsetEpoch=Optional[2], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=2}}
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-6 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=7, clientId=consumer-low-priority-email-group-7, groupId=low-priority-email-group] Setting offset for partition low-priority-email-7 to the committed offset FetchPosition{offset=13303, offsetEpoch=Optional[2], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=2}}
2023-11-08T02:32:21.923-0500 consumer-low-priority-email-pool-4 org.apache.kafka.clients.consumer.internals.ConsumerCoordinator INFO: [Consumer instanceId=5, clientId=consumer-low-priority-email-group-5, groupId=low-priority-email-group] Setting offset for partition low-priority-email-5 to the committed offset FetchPosition{offset=13338, offsetEpoch=Optional[2], currentLeader=LeaderAndEpoch{leader=Optional[192.168.100.60:9093 (id: 1 rack: null)], epoch=2}}
The full logs can be downloaded here without CSDN points.
References
Closing
This largely completes our Kafka messaging system. The key code has been explained across the articles. To understand the full design and implementation, get the complete code from the source repository.
The next Kafka article is planned to explain synchronizing Oracle data to PostgreSQL with Kafka Connect.