From 34f4ad6a8be6aa5145f14cfdbcb643b5070cdc5b Mon Sep 17 00:00:00 2001 From: abbccdda Date: Tue, 6 Aug 2019 13:08:39 -0700 Subject: [PATCH 1/6] condensed changes --- .../kafka/clients/producer/KafkaProducer.java | 22 ++++++ .../kafka/clients/producer/Producer.java | 7 ++ .../internals/TransactionManager.java | 47 +++++++++++-- .../internals/TransactionManagerTest.java | 16 ++--- .../apache/kafka/streams/StreamsConfig.java | 28 ++++++++ .../internals/AssignedStandbyTasks.java | 2 +- .../internals/AssignedStreamsTasks.java | 48 +++---------- .../processor/internals/AssignedTasks.java | 53 ++++++++++++-- .../processor/internals/StreamTask.java | 55 ++++++++++----- .../processor/internals/StreamThread.java | 69 ++++++++++++++++--- .../streams/processor/internals/Task.java | 15 ++++ .../processor/internals/StreamTaskTest.java | 20 +++--- .../processor/internals/StreamThreadTest.java | 36 +++++----- .../StreamThreadStateStoreProviderTest.java | 2 +- .../kafka/streams/TopologyTestDriver.java | 2 +- 15 files changed, 306 insertions(+), 116 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java index c586fb41be414..09b4201152216 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java @@ -21,6 +21,8 @@ import org.apache.kafka.clients.ClientUtils; import org.apache.kafka.clients.KafkaClient; import org.apache.kafka.clients.NetworkClient; +import org.apache.kafka.clients.consumer.Consumer; +import org.apache.kafka.clients.consumer.ConsumerGroupMetadata; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.clients.consumer.OffsetCommitCallback; @@ -78,6 +80,7 @@ import java.util.Map; import java.util.Objects; import java.util.Properties; +import java.util.UUID; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; @@ -255,6 +258,7 @@ public class KafkaProducer implements Producer { private final ProducerInterceptors interceptors; private final ApiVersions apiVersions; private final TransactionManager transactionManager; + private ConsumerGroupMetadata consumerGroupMetadata; /** * A producer is instantiated by providing a set of key-value pairs as configuration. Valid configuration strings @@ -618,11 +622,20 @@ private static int parseAcks(String acksString) { public void initTransactions() { throwIfNoTransactionManager(); throwIfProducerClosed(); + maybeAllocateTransactionalId(); + TransactionalRequestResult result = transactionManager.initializeTransactions(); sender.wakeup(); result.await(maxBlockTimeMs, TimeUnit.MILLISECONDS); } + private void maybeAllocateTransactionalId() { + String allocatedTransactionalId = transactionManager.transactionalId(); + if (allocatedTransactionalId == null || allocatedTransactionalId.isEmpty()) { + transactionManager.setTransactionalId("thread-producer-" + UUID.randomUUID().toString()); + } + } + /** * Should be called before the start of each new transaction. Note that prior to the first invocation * of this method, you must invoke {@link #initTransactions()} exactly one time. @@ -668,6 +681,15 @@ public void beginTransaction() throws ProducerFencedException { */ public void sendOffsetsToTransaction(Map offsets, String consumerGroupId) throws ProducerFencedException { + sendOffsetToTransactionInternal(offsets, consumerGroupId); + } + + public void sendOffsetsToTransaction(Map offsets) throws ProducerFencedException { + sendOffsetToTransactionInternal(offsets, consumerGroupMetadata.groupId()); + } + + private void sendOffsetToTransactionInternal(Map offsets, + String consumerGroupId) { throwIfNoTransactionManager(); throwIfProducerClosed(); TransactionalRequestResult result = transactionManager.sendOffsetsToTransaction(offsets, consumerGroupId); diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/Producer.java b/clients/src/main/java/org/apache/kafka/clients/producer/Producer.java index 96d487a79235b..60a326ba61fe4 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/Producer.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/Producer.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.clients.producer; +import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.common.Metric; import org.apache.kafka.common.MetricName; import org.apache.kafka.common.PartitionInfo; @@ -53,6 +54,12 @@ public interface Producer extends Closeable { void sendOffsetsToTransaction(Map offsets, String consumerGroupId) throws ProducerFencedException; + + /** + * See {@link KafkaProducer#sendOffsetsToTransaction(Map)} + */ + void sendOffsetsToTransaction(Map offsets) throws ProducerFencedException; + /** * See {@link KafkaProducer#commitTransaction()} */ diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java index cd091542dc8db..14919d054fb30 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java @@ -19,6 +19,7 @@ import org.apache.kafka.clients.ClientResponse; import org.apache.kafka.clients.RequestCompletionHandler; import org.apache.kafka.clients.consumer.OffsetAndMetadata; +import org.apache.kafka.clients.consumer.internals.ConsumerGroupMetadata; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; @@ -32,6 +33,7 @@ import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.message.FindCoordinatorRequestData; import org.apache.kafka.common.message.InitProducerIdRequestData; +import org.apache.kafka.common.message.TxnOffsetCommitRequestData; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.record.DefaultRecordBatch; import org.apache.kafka.common.record.RecordBatch; @@ -83,7 +85,7 @@ public class TransactionManager { private static final int NO_LAST_ACKED_SEQUENCE_NUMBER = -1; private final Logger log; - private final String transactionalId; + private String transactionalId; private final int transactionTimeoutMs; private static class TopicPartitionBookkeeper { @@ -200,6 +202,7 @@ public void resetSequenceNumbers(Consumer resetSequence) { private volatile RuntimeException lastError = null; private volatile ProducerIdAndEpoch producerIdAndEpoch; private volatile boolean transactionStarted = false; + private org.apache.kafka.clients.consumer.Consumer consumer; private enum State { UNINITIALIZED, @@ -272,6 +275,14 @@ public TransactionManager(LogContext logContext, String transactionalId, int tra this(new LogContext(), null, 0, 100L); } + public void setTransactionalId(String transactionalId) { + this.transactionalId = transactionalId; + } + + public void setConsumer(org.apache.kafka.clients.consumer.Consumer consumer) { + this.consumer = consumer; + } + public synchronized TransactionalRequestResult initializeTransactions() { return handleCachedTransactionRequestResult(() -> { transitionTo(State.INITIALIZING); @@ -986,8 +997,20 @@ private TxnOffsetCommitHandler txnOffsetCommitHandler(TransactionalRequestResult offsetAndMetadata.metadata(), offsetAndMetadata.leaderEpoch()); pendingTxnOffsetCommits.put(entry.getKey(), committedOffset); } - TxnOffsetCommitRequest.Builder builder = new TxnOffsetCommitRequest.Builder(transactionalId, consumerGroupId, - producerIdAndEpoch.producerId, producerIdAndEpoch.epoch, pendingTxnOffsetCommits); + + + ConsumerGroupMetadata metadata = (ConsumerGroupMetadata) consumer.groupMetadata(); + TxnOffsetCommitRequest.Builder builder = new TxnOffsetCommitRequest.Builder( + new TxnOffsetCommitRequestData() + .setTransactionalId(transactionalId) + .setGroupId(consumerGroupId) + .setProducerId(producerIdAndEpoch.producerId) + .setProducerEpoch(producerIdAndEpoch.epoch) + .setTopics(TxnOffsetCommitRequest.getTopics(pendingTxnOffsetCommits)) + .setGenerationId(metadata.generation()) + .setMemberId(metadata.memberId()) + .setGroupInstanceId(metadata.groupInstanceId().orElse(null)) + ); return new TxnOffsetCommitHandler(result, builder); } @@ -1194,6 +1217,16 @@ public void handleResponse(AbstractResponse response) { } else if (error == Errors.INVALID_PRODUCER_EPOCH) { fatalError(error.exception()); return; + } else if (error == Errors.ILLEGAL_GENERATION) { + fatalError(error.exception()); + return; + } else if (error == Errors.FENCED_INSTANCE_ID) { + fatalError(error.exception()); + return; + } else if (error == Errors.UNKNOWN_MEMBER_ID) { + // What should we do? The consumer internal is not available here. + reenqueue(); + return; } else if (error == Errors.TRANSACTIONAL_ID_AUTHORIZATION_FAILED) { fatalError(error.exception()); return; @@ -1432,7 +1465,7 @@ FindCoordinatorRequest.CoordinatorType coordinatorType() { @Override String coordinatorKey() { - return builder.consumerGroupId(); + return builder.data.groupId(); } @Override @@ -1441,7 +1474,7 @@ public void handleResponse(AbstractResponse response) { boolean coordinatorReloaded = false; Map errors = txnOffsetCommitResponse.errors(); - log.debug("Received TxnOffsetCommit response for consumer group {}: {}", builder.consumerGroupId(), + log.debug("Received TxnOffsetCommit response for consumer group {}: {}", builder.data.groupId(), errors); for (Map.Entry entry : errors.entrySet()) { @@ -1454,14 +1487,14 @@ public void handleResponse(AbstractResponse response) { || error == Errors.REQUEST_TIMED_OUT) { if (!coordinatorReloaded) { coordinatorReloaded = true; - lookupCoordinator(FindCoordinatorRequest.CoordinatorType.GROUP, builder.consumerGroupId()); + lookupCoordinator(FindCoordinatorRequest.CoordinatorType.GROUP, builder.data.groupId()); } } else if (error == Errors.UNKNOWN_TOPIC_OR_PARTITION || error == Errors.COORDINATOR_LOAD_IN_PROGRESS) { // If the topic is unknown or the coordinator is loading, retry with the current coordinator continue; } else if (error == Errors.GROUP_AUTHORIZATION_FAILED) { - abortableError(GroupAuthorizationException.forGroupId(builder.consumerGroupId())); + abortableError(GroupAuthorizationException.forGroupId(builder.data.groupId())); break; } else if (error == Errors.TRANSACTIONAL_ID_AUTHORIZATION_FAILED || error == Errors.INVALID_PRODUCER_EPOCH diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java index cca5771002cd8..a885486062cf3 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java @@ -919,7 +919,7 @@ public void testUnsupportedForMessageFormatInTxnOffsetCommit() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.NONE, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled, TxnOffsetCommit Enqueued @@ -1089,7 +1089,7 @@ public void testGroupAuthorizationFailureInFindCoordinator() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(new TopicPartition("foo", 0), new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(new TopicPartition("foo", 0), new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.NONE, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled, TxnOffsetCommit Enqueued @@ -1121,7 +1121,7 @@ public void testGroupAuthorizationFailureInTxnOffsetCommit() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(tp1, new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(tp1, new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.NONE, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled, TxnOffsetCommit Enqueued @@ -1157,7 +1157,7 @@ public void testTransactionalIdAuthorizationFailureInAddOffsetsToTxn() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.TRANSACTIONAL_ID_AUTHORIZATION_FAILED, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled @@ -1182,7 +1182,7 @@ public void testTransactionalIdAuthorizationFailureInTxnOffsetCommit() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.NONE, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled, TxnOffsetCommit Enqueued @@ -2826,9 +2826,9 @@ private void prepareTxnOffsetCommitResponse(final String consumerGroupId, Map txnOffsetCommitResponse) { client.prepareResponse(request -> { TxnOffsetCommitRequest txnOffsetCommitRequest = (TxnOffsetCommitRequest) request; - assertEquals(consumerGroupId, txnOffsetCommitRequest.consumerGroupId()); - assertEquals(producerId, txnOffsetCommitRequest.producerId()); - assertEquals(producerEpoch, txnOffsetCommitRequest.producerEpoch()); + assertEquals(consumerGroupId, txnOffsetCommitRequest.data.groupId()); + assertEquals(producerId, txnOffsetCommitRequest.data.producerId()); + assertEquals(producerEpoch, txnOffsetCommitRequest.data.producerEpoch()); return true; }, new TxnOffsetCommitResponse(0, txnOffsetCommitResponse)); } diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java index f08eecaac4743..66a4f4132be2c 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java @@ -251,6 +251,34 @@ public class StreamsConfig extends AbstractConfig { @SuppressWarnings("WeakerAccess") public static final String UPGRADE_FROM_11 = "1.1"; + /** + * Config value for parameter {@link #UPGRADE_FROM_CONFIG "upgrade.from"} for upgrading an application from version {@code 2.0.x}. + */ + @SuppressWarnings("WeakerAccess") + public static final String UPGRADE_FROM_20 = "2.0"; + + + /** + * Config value for parameter {@link #UPGRADE_FROM_CONFIG "upgrade.from"} for upgrading an application from version {@code 2.1.x}. + */ + @SuppressWarnings("WeakerAccess") + public static final String UPGRADE_FROM_21 = "2.1"; + + + + /** + * Config value for parameter {@link #UPGRADE_FROM_CONFIG "upgrade.from"} for upgrading an application from version {@code 2.2.x}. + */ + @SuppressWarnings("WeakerAccess") + public static final String UPGRADE_FROM_22 = "2.2"; + + + /** + * Config value for parameter {@link #UPGRADE_FROM_CONFIG "upgrade.from"} for upgrading an application from version {@code 2.3.x}. + */ + @SuppressWarnings("WeakerAccess") + public static final String UPGRADE_FROM_23 = "2.3"; + /** * Config value for parameter {@link #PROCESSING_GUARANTEE_CONFIG "processing.guarantee"} for at-least-once processing guarantees. */ diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStandbyTasks.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStandbyTasks.java index a99e45147b9ec..30e4914fc98a8 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStandbyTasks.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStandbyTasks.java @@ -21,7 +21,7 @@ class AssignedStandbyTasks extends AssignedTasks { AssignedStandbyTasks(final LogContext logContext) { - super(logContext, "standby task"); + super(logContext, "standby task", null, null); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasks.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasks.java index f0019ec96722c..cfd37d1af45d2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasks.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasks.java @@ -16,9 +16,11 @@ */ package org.apache.kafka.streams.processor.internals; +import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.utils.LogContext; +import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.errors.TaskMigratedException; import org.apache.kafka.streams.processor.TaskId; @@ -35,9 +37,13 @@ class AssignedStreamsTasks extends AssignedTasks implements Restorin private final Map restoring = new HashMap<>(); private final Set restoredPartitions = new HashSet<>(); private final Map restoringByPartition = new HashMap<>(); + private final Producer threadProducer; - AssignedStreamsTasks(final LogContext logContext) { - super(logContext, "stream task"); + AssignedStreamsTasks(final LogContext logContext, + final Producer threadProducer, + final Time time) { + super(logContext, "stream task", threadProducer, time); + this.threadProducer = threadProducer; } @Override @@ -136,41 +142,9 @@ void addToRestoring(final StreamTask task) { * or if the task producer got fenced (EOS) */ int maybeCommitPerUserRequested() { - int committed = 0; - RuntimeException firstException = null; - - for (final Iterator it = running().iterator(); it.hasNext(); ) { - final StreamTask task = it.next(); - try { - if (task.commitRequested() && task.commitNeeded()) { - task.commit(); - committed++; - log.debug("Committed active task {} per user request in", task.id()); - } - } catch (final TaskMigratedException e) { - log.info("Failed to commit {} since it got migrated to another thread already. " + - "Closing it as zombie before triggering a new rebalance.", task.id()); - final RuntimeException fatalException = closeZombieTask(task); - if (fatalException != null) { - throw fatalException; - } - it.remove(); - throw e; - } catch (final RuntimeException t) { - log.error("Failed to commit StreamTask {} due to the following error:", - task.id(), - t); - if (firstException == null) { - firstException = t; - } - } - } - - if (firstException != null) { - throw firstException; - } - - return committed; + return commitInternal(log, + threadProducer, + task -> task.commitRequested() && task.commitNeeded()); } /** diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java index e1bfe37247267..132520f9c1b99 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java @@ -16,8 +16,11 @@ */ package org.apache.kafka.streams.processor.internals; +import org.apache.kafka.clients.consumer.OffsetAndMetadata; +import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.utils.LogContext; +import org.apache.kafka.common.utils.Time; import org.apache.kafka.streams.errors.LockException; import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.errors.TaskMigratedException; @@ -42,15 +45,21 @@ abstract class AssignedTasks { private final Map created = new HashMap<>(); private final Map suspended = new HashMap<>(); private final Set previousActiveTasks = new HashSet<>(); + private final Producer eosProducer; + private final Time time; // IQ may access this map. final Map running = new ConcurrentHashMap<>(); private final Map runningByPartition = new HashMap<>(); AssignedTasks(final LogContext logContext, - final String taskTypeName) { + final String taskTypeName, + final Producer eosProducer, + final Time time) { this.taskTypeName = taskTypeName; this.log = logContext.logger(getClass()); + this.eosProducer = eosProducer; + this.time = time; } void addNewTask(final T task) { @@ -276,13 +285,36 @@ Set previousTaskIds() { * or if the task producer got fenced (EOS) */ int commit() { + return commitInternal( + log, + eosProducer, + Task::commitNeeded + ); + } + + public interface TaskStatus { + boolean needsCommit(Task task); + } + + protected int commitInternal(Logger log, + Producer eosProducer, + TaskStatus taskStatus) { int committed = 0; RuntimeException firstException = null; + + Map pendingOffsets = new HashMap<>(); + List tasks = new ArrayList<>(); + for (final Iterator it = running().iterator(); it.hasNext(); ) { final T task = it.next(); try { - if (task.commitNeeded()) { - task.commit(); + if (taskStatus.needsCommit(task)) { + if (eosProducer != null) { + pendingOffsets.putAll(task.getPendingOffsets()); + tasks.add(task); + } else { + task.commit(); + } committed++; } } catch (final TaskMigratedException e) { @@ -296,9 +328,7 @@ int commit() { throw e; } catch (final RuntimeException t) { log.error("Failed to commit {} {} due to the following error:", - taskTypeName, - task.id(), - t); + taskTypeName, task.id(), t); if (firstException == null) { firstException = t; } @@ -309,6 +339,17 @@ int commit() { throw firstException; } + if (!pendingOffsets.isEmpty()) { + final long startNs = time.nanoseconds(); + eosProducer.sendOffsetsToTransaction(pendingOffsets); + eosProducer.commitTransaction(); + long commitLatency = time.nanoseconds() - startNs; + for (Task task : tasks) { + task.markCommitDone(commitLatency); + } + eosProducer.beginTransaction(); + } + return committed; } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index 836330f242617..562500f9d4e75 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -82,6 +82,7 @@ public class StreamTask extends AbstractTask implements ProcessorNodePunctuator private Sensor closeTaskSensor; private long idleStartTime; private Producer producer; + private final boolean isTaskProducer; private boolean commitRequested = false; private boolean transactionInFlight = false; @@ -151,8 +152,9 @@ public StreamTask(final TaskId id, final StateDirectory stateDirectory, final ThreadCache cache, final Time time, - final ProducerSupplier producerSupplier) { - this(id, partitions, topology, consumer, changelogReader, config, metrics, stateDirectory, cache, time, producerSupplier, null); + final ProducerSupplier producerSupplier, + boolean isTaskProducer) { + this(id, partitions, topology, consumer, changelogReader, config, metrics, stateDirectory, cache, time, producerSupplier, null, isTaskProducer); } public StreamTask(final TaskId id, @@ -166,9 +168,11 @@ public StreamTask(final TaskId id, final ThreadCache cache, final Time time, final ProducerSupplier producerSupplier, - final RecordCollector recordCollector) { + final RecordCollector recordCollector, + boolean isTaskProducer) { super(id, partitions, topology, consumer, changelogReader, false, stateDirectory, config); + this.isTaskProducer = isTaskProducer; this.time = time; this.producerSupplier = producerSupplier; this.producer = producerSupplier.get(); @@ -228,7 +232,7 @@ public StreamTask(final TaskId id, // initialize transactions if eos is turned on, which will block if the previous transaction has not // completed yet; do not start the first transaction until the topology has been initialized later - if (eosEnabled) { + if (isTaskProducer) { initializeTransactions(); } } @@ -252,7 +256,7 @@ public boolean initializeStateStores() { public void initializeTopology() { initTopology(); - if (eosEnabled) { + if (isTaskProducer) { try { this.producer.beginTransaction(); } catch (final ProducerFencedException fatal) { @@ -278,7 +282,7 @@ public void initializeTopology() { @Override public void resume() { log.debug("Resuming"); - if (eosEnabled) { + if (isTaskProducer) { if (producer != null) { throw new IllegalStateException("Task producer should be null."); } @@ -455,16 +459,10 @@ void commit(final boolean startNewTransaction) { stateMgr.checkpoint(activeTaskCheckpointableOffsets()); } - final Map consumedOffsetsAndMetadata = new HashMap<>(consumedOffsets.size()); - for (final Map.Entry entry : consumedOffsets.entrySet()) { - final TopicPartition partition = entry.getKey(); - final long offset = entry.getValue() + 1; - consumedOffsetsAndMetadata.put(partition, new OffsetAndMetadata(offset)); - stateMgr.putOffsetLimit(partition, offset); - } + final Map consumedOffsetsAndMetadata = getPendingOffsets(); try { - if (eosEnabled) { + if (isTaskProducer) { producer.sendOffsetsToTransaction(consumedOffsetsAndMetadata, applicationId); producer.commitTransaction(); transactionInFlight = false; @@ -479,9 +477,26 @@ void commit(final boolean startNewTransaction) { throw new TaskMigratedException(this, error); } + markCommitDone(time.nanoseconds() - startNs); + } + + @Override + public Map getPendingOffsets() { + final Map consumedOffsetsAndMetadata = new HashMap<>(consumedOffsets.size()); + for (final Map.Entry entry : consumedOffsets.entrySet()) { + final TopicPartition partition = entry.getKey(); + final long offset = entry.getValue() + 1; + consumedOffsetsAndMetadata.put(partition, new OffsetAndMetadata(offset)); + stateMgr.putOffsetLimit(partition, offset); + } + return consumedOffsetsAndMetadata; + } + + @Override + public void markCommitDone(long commitLatency) { commitNeeded = false; commitRequested = false; - taskMetrics.taskCommitTimeSensor.record(time.nanoseconds() - startNs); + taskMetrics.taskCommitTimeSensor.record(commitLatency); } @Override @@ -598,7 +613,7 @@ void suspend(final boolean clean, } private void maybeAbortTransactionAndCloseRecordCollector(final boolean isZombie) { - if (eosEnabled && !isZombie) { + if (isTaskProducer && !isZombie) { try { if (transactionInFlight) { producer.abortTransaction(); @@ -616,7 +631,7 @@ private void maybeAbortTransactionAndCloseRecordCollector(final boolean isZombie } } - if (eosEnabled) { + if (isTaskProducer) { try { recordCollector.close(); } catch (final Throwable e) { @@ -844,7 +859,8 @@ void requestCommit() { /** * Whether or not a request has been made to commit the current state */ - boolean commitRequested() { + @Override + public boolean commitRequested() { return commitRequested; } @@ -859,7 +875,8 @@ Producer getProducer() { private void initializeTransactions() { try { - producer.initTransactions(); + producer.initTransactions(consumer); + producer.beginTransaction(); } catch (final TimeoutException retriable) { log.error( "Timeout exception caught when initializing transactions for task {}. " + diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index d3efa9e0e05de..1a6cdc1a9ed95 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -29,6 +29,7 @@ import org.apache.kafka.common.Metric; import org.apache.kafka.common.MetricName; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.TimeoutException; import org.apache.kafka.common.metrics.Metrics; import org.apache.kafka.common.metrics.Sensor; import org.apache.kafka.common.serialization.ByteArrayDeserializer; @@ -61,6 +62,7 @@ import java.util.UUID; import java.util.concurrent.atomic.AtomicInteger; +import static java.lang.String.format; import static java.util.Collections.singleton; public class StreamThread extends Thread { @@ -451,7 +453,8 @@ StreamTask createTask(final Consumer consumer, stateDirectory, cache, time, - () -> createProducer(taskId)); + () -> createProducer(taskId), + threadProducer == null); } private Producer createProducer(final TaskId id) { @@ -563,6 +566,8 @@ StandbyTask createTask(final Consumer consumer, final Consumer consumer; final InternalTopologyBuilder builder; + final boolean eosEnabled; + public static StreamThread create(final InternalTopologyBuilder builder, final StreamsConfig config, final KafkaClientSupplier clientSupplier, @@ -590,7 +595,9 @@ public static StreamThread create(final InternalTopologyBuilder builder, Producer threadProducer = null; final boolean eosEnabled = StreamsConfig.EXACTLY_ONCE.equals(config.getString(StreamsConfig.PROCESSING_GUARANTEE_CONFIG)); - if (!eosEnabled) { + final String upgradeFrom = config.getString(StreamsConfig.UPGRADE_FROM_CONFIG); + + if (useThreadProducer(eosEnabled, upgradeFrom)) { final Map producerConfigs = config.getProducerConfigs(getThreadProducerClientId(threadClientId)); log.info("Creating shared producer client"); threadProducer = clientSupplier.getProducer(producerConfigs); @@ -629,7 +636,7 @@ public static StreamThread create(final InternalTopologyBuilder builder, activeTaskCreator, standbyTaskCreator, adminClient, - new AssignedStreamsTasks(logContext), + new AssignedStreamsTasks(logContext, threadProducer, time), new AssignedStandbyTasks(logContext)); log.info("Creating consumer client"); @@ -644,25 +651,46 @@ public static StreamThread create(final InternalTopologyBuilder builder, consumerConfigs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "none"); } - final Consumer consumer = clientSupplier.getConsumer(consumerConfigs); - taskManager.setConsumer(consumer); + final Consumer mainConsumer = clientSupplier.getConsumer(consumerConfigs); + taskManager.setConsumer(mainConsumer); return new StreamThread( time, config, threadProducer, restoreConsumer, - consumer, + mainConsumer, originalReset, taskManager, streamsMetrics, builder, threadClientId, logContext, - assignmentErrorCode) + assignmentErrorCode, + eosEnabled) .updateThreadMetadata(getSharedAdminClientId(clientId)); } + private static boolean useThreadProducer(final boolean eosEnabled, final String upgradeFrom) { + switch (upgradeFrom) { + case StreamsConfig.UPGRADE_FROM_0100: + case StreamsConfig.UPGRADE_FROM_0101: + case StreamsConfig.UPGRADE_FROM_0102: + case StreamsConfig.UPGRADE_FROM_0110: + case StreamsConfig.UPGRADE_FROM_10: + case StreamsConfig.UPGRADE_FROM_11: + case StreamsConfig.UPGRADE_FROM_20: + case StreamsConfig.UPGRADE_FROM_21: + case StreamsConfig.UPGRADE_FROM_22: + case StreamsConfig.UPGRADE_FROM_23: + // If upgrading from an older version, we shall continue using the eos producer setup. + return !eosEnabled; + default: + // If not upgrading, start from 2.4 there will be no task level producer setup. + return true; + } + } + public StreamThread(final Time time, final StreamsConfig config, final Producer producer, @@ -674,7 +702,8 @@ public StreamThread(final Time time, final InternalTopologyBuilder builder, final String threadClientId, final LogContext logContext, - final AtomicInteger assignmentErrorCode) { + final AtomicInteger assignmentErrorCode, + boolean eosEnabled) { super(threadClientId); this.stateLock = new Object(); @@ -715,6 +744,8 @@ public StreamThread(final Time time, this.commitTimeMs = config.getLong(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG); this.numIterations = 1; + + this.eosEnabled = eosEnabled; } private static final class InternalConsumerConfig extends ConsumerConfig { @@ -759,6 +790,7 @@ public void run() { } boolean cleanRun = false; try { + maybeInitializeTransactions(); runLoop(); cleanRun = true; } catch (final KafkaException e) { @@ -775,6 +807,27 @@ public void run() { } } + private void maybeInitializeTransactions() { + try { + // This is a thread-level txn producer + if (eosEnabled && producer != null) { + producer.initTransactions(); + } + } catch (final TimeoutException retriable) { + log.error( + "Timeout exception caught when initializing transactions for current stream thread. " + + "This might happen if the broker is slow to respond, if the network connection to " + + "the broker was interrupted, or if similar circumstances arise. " + + "You can increase producer parameter `max.block.ms` to increase this timeout.", + retriable + ); + throw new StreamsException( + format("%sFailed to initialize stream thread due to timeout.", logPrefix), + retriable + ); + } + } + private void setRebalanceException(final Throwable rebalanceException) { this.rebalanceException = rebalanceException; } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java index 812e7e1131ca2..af09fed9c4563 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java @@ -16,13 +16,16 @@ */ package org.apache.kafka.streams.processor.internals; +import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.streams.errors.StreamsException; import org.apache.kafka.streams.processor.ProcessorContext; import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.processor.TaskId; +import sun.reflect.generics.reflectiveObjects.NotImplementedException; import java.util.Collection; +import java.util.Map; import java.util.Set; public interface Task { @@ -36,10 +39,22 @@ public interface Task { boolean commitNeeded(); + default boolean commitRequested() { + return false; + } + void initializeTopology(); void commit(); + default Map getPendingOffsets() { + throw new NotImplementedException(); + } + + default void markCommitDone(long commitLatency) { + throw new NotImplementedException(); + } + void suspend(); void resume(); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java index a70c10c3a7216..b87ab08bd2354 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java @@ -246,8 +246,8 @@ public void initTransactions() { throw new TimeoutException("test"); } }, - null - ); + null, + false); fail("Expected an exception"); } catch (final StreamsException expected) { // make sure we log the explanation as an ERROR @@ -300,8 +300,8 @@ public void initTransactions() { } } }, - null - ); + null, + false); testTask.initializeTopology(); testTask.suspend(); timeOut.set(true); @@ -849,7 +849,7 @@ public void shouldFlushRecordCollectorOnFlushState() { public void flush() { flushed.set(true); } - }); + }, false); streamTask.flushState(); assertTrue(flushed.get()); } @@ -1424,7 +1424,7 @@ public void shouldReturnOffsetsForRepartitionTopicsForPurging() { stateDirectory, null, time, - () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer)); + () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer), false); task.initializeStateStores(); task.initializeTopology(); @@ -1496,7 +1496,7 @@ private StreamTask createStatefulTask(final StreamsConfig config, final boolean stateDirectory, null, time, - () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer)); + () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer), false); } private StreamTask createStatefulTaskThatThrowsExceptionOnClose() { @@ -1517,7 +1517,7 @@ private StreamTask createStatefulTaskThatThrowsExceptionOnClose() { stateDirectory, null, time, - () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer)); + () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer), false); } private StreamTask createStatelessTask(final StreamsConfig streamsConfig) { @@ -1542,7 +1542,7 @@ private StreamTask createStatelessTask(final StreamsConfig streamsConfig) { stateDirectory, null, time, - () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer)); + () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer), false); } // this task will throw exception when processing (on partition2), flushing, suspending and closing @@ -1568,7 +1568,7 @@ private StreamTask createTaskThatThrowsException(final boolean enableEos) { stateDirectory, null, time, - () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer)) { + () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer), false) { @Override protected void flushState() { throw new RuntimeException("KABOOM!"); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java index aff5d6c74861c..5e768b1000b45 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java @@ -353,8 +353,8 @@ public void shouldNotCommitBeforeTheCommitInterval() { internalTopologyBuilder, clientId, new LogContext(""), - new AtomicInteger() - ); + new AtomicInteger(), + false); thread.setNow(mockTime.milliseconds()); thread.maybeCommit(); mockTime.sleep(commitInterval - 10L); @@ -479,8 +479,8 @@ public void shouldNotCauseExceptionIfNothingCommitted() { internalTopologyBuilder, clientId, new LogContext(""), - new AtomicInteger() - ); + new AtomicInteger(), + false); thread.setNow(mockTime.milliseconds()); thread.maybeCommit(); mockTime.sleep(commitInterval - 10L); @@ -514,8 +514,8 @@ public void shouldCommitAfterTheCommitInterval() { internalTopologyBuilder, clientId, new LogContext(""), - new AtomicInteger() - ); + new AtomicInteger(), + false); thread.setNow(mockTime.milliseconds()); thread.maybeCommit(); mockTime.sleep(commitInterval + 1); @@ -664,8 +664,8 @@ public void shouldShutdownTaskManagerOnClose() { internalTopologyBuilder, clientId, new LogContext(""), - new AtomicInteger() - ).updateThreadMetadata(getSharedAdminClientId(clientId)); + new AtomicInteger(), + false).updateThreadMetadata(getSharedAdminClientId(clientId)); thread.setStateListener( (t, newState, oldState) -> { if (oldState == StreamThread.State.CREATED && newState == StreamThread.State.STARTING) { @@ -697,8 +697,8 @@ public void shouldShutdownTaskManagerOnCloseWithoutStart() { internalTopologyBuilder, clientId, new LogContext(""), - new AtomicInteger() - ).updateThreadMetadata(getSharedAdminClientId(clientId)); + new AtomicInteger(), + false).updateThreadMetadata(getSharedAdminClientId(clientId)); thread.shutdown(); EasyMock.verify(taskManager); } @@ -769,8 +769,8 @@ private void setStreamThread(final StreamThread streamThread) { internalTopologyBuilder, clientId, new LogContext(""), - new AtomicInteger() - ).updateThreadMetadata(getSharedAdminClientId(clientId)); + new AtomicInteger(), + false).updateThreadMetadata(getSharedAdminClientId(clientId)); mockStreamThreadConsumer.setStreamThread(thread); mockStreamThreadConsumer.assign(assignedPartitions); @@ -803,8 +803,8 @@ public void shouldOnlyShutdownOnce() { internalTopologyBuilder, clientId, new LogContext(""), - new AtomicInteger() - ).updateThreadMetadata(getSharedAdminClientId(clientId)); + new AtomicInteger(), + false).updateThreadMetadata(getSharedAdminClientId(clientId)); thread.shutdown(); // Execute the run method. Verification of the mock will check that shutdown was only done once thread.run(); @@ -1634,8 +1634,8 @@ public void producerMetricsVerificationWithoutEOS() { internalTopologyBuilder, clientId, new LogContext(""), - new AtomicInteger() - ); + new AtomicInteger(), + false); final MetricName testMetricName = new MetricName("test_metric", "", "", new HashMap<>()); final Metric testMetric = new KafkaMetric( new Object(), @@ -1673,8 +1673,8 @@ public void adminClientMetricsVerification() { internalTopologyBuilder, clientId, new LogContext(""), - new AtomicInteger() - ); + new AtomicInteger(), + false); final MetricName testMetricName = new MetricName("test_metric", "", "", new HashMap<>()); final Metric testMetric = new KafkaMetric( new Object(), diff --git a/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java b/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java index f48c31c1d91c2..5780ddd6c19e3 100644 --- a/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/state/internals/StreamThreadStateStoreProviderTest.java @@ -317,7 +317,7 @@ private StreamTask createStreamsTask(final StreamsConfig streamsConfig, stateDirectory, null, new MockTime(), - () -> clientSupplier.getProducer(new HashMap<>())) { + () -> clientSupplier.getProducer(new HashMap<>()), false) { @Override protected void updateOffsetLimits() {} }; diff --git a/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java b/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java index 2b1428f9f8968..5370c3906caf9 100644 --- a/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java +++ b/streams/test-utils/src/main/java/org/apache/kafka/streams/TopologyTestDriver.java @@ -374,7 +374,7 @@ public void onRestoreEnd(final TopicPartition topicPartition, final String store stateDirectory, cache, mockWallClockTime, - () -> producer); + () -> producer, false); task.initializeStateStores(); task.initializeTopology(); ((InternalProcessorContext) task.context()).setRecordContext(new ProcessorRecordContext( From dbafff89d9dd914ec4a28ec352d1ff9868edc5be Mon Sep 17 00:00:00 2001 From: abbccdda Date: Tue, 6 Aug 2019 13:36:12 -0700 Subject: [PATCH 2/6] add producer --- .../kafka/clients/producer/KafkaProducer.java | 16 ++------ .../kafka/clients/producer/Producer.java | 7 ---- .../internals/TransactionManager.java | 40 +++---------------- .../internals/TransactionManagerTest.java | 16 ++++---- .../internals/AssignedStandbyTasks.java | 2 +- .../internals/AssignedStreamsTasks.java | 10 ++--- .../processor/internals/AssignedTasks.java | 21 +++++----- .../processor/internals/StreamTask.java | 8 ++-- .../processor/internals/StreamThread.java | 14 +++---- .../streams/processor/internals/Task.java | 6 +-- .../internals/AssignedStreamsTasksTest.java | 3 +- .../processor/internals/StreamThreadTest.java | 2 +- 12 files changed, 52 insertions(+), 93 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java index 09b4201152216..837fe29607c47 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java @@ -21,8 +21,6 @@ import org.apache.kafka.clients.ClientUtils; import org.apache.kafka.clients.KafkaClient; import org.apache.kafka.clients.NetworkClient; -import org.apache.kafka.clients.consumer.Consumer; -import org.apache.kafka.clients.consumer.ConsumerGroupMetadata; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndMetadata; import org.apache.kafka.clients.consumer.OffsetCommitCallback; @@ -258,7 +256,6 @@ public class KafkaProducer implements Producer { private final ProducerInterceptors interceptors; private final ApiVersions apiVersions; private final TransactionManager transactionManager; - private ConsumerGroupMetadata consumerGroupMetadata; /** * A producer is instantiated by providing a set of key-value pairs as configuration. Valid configuration strings @@ -632,7 +629,9 @@ public void initTransactions() { private void maybeAllocateTransactionalId() { String allocatedTransactionalId = transactionManager.transactionalId(); if (allocatedTransactionalId == null || allocatedTransactionalId.isEmpty()) { - transactionManager.setTransactionalId("thread-producer-" + UUID.randomUUID().toString()); + String threadProducerId = "thread-producer-" + UUID.randomUUID().toString(); + log.info("Allocating thread producer id: {}", threadProducerId); + transactionManager.setTransactionalId(threadProducerId); } } @@ -681,15 +680,6 @@ public void beginTransaction() throws ProducerFencedException { */ public void sendOffsetsToTransaction(Map offsets, String consumerGroupId) throws ProducerFencedException { - sendOffsetToTransactionInternal(offsets, consumerGroupId); - } - - public void sendOffsetsToTransaction(Map offsets) throws ProducerFencedException { - sendOffsetToTransactionInternal(offsets, consumerGroupMetadata.groupId()); - } - - private void sendOffsetToTransactionInternal(Map offsets, - String consumerGroupId) { throwIfNoTransactionManager(); throwIfProducerClosed(); TransactionalRequestResult result = transactionManager.sendOffsetsToTransaction(offsets, consumerGroupId); diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/Producer.java b/clients/src/main/java/org/apache/kafka/clients/producer/Producer.java index 60a326ba61fe4..96d487a79235b 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/Producer.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/Producer.java @@ -16,7 +16,6 @@ */ package org.apache.kafka.clients.producer; -import org.apache.kafka.clients.consumer.Consumer; import org.apache.kafka.common.Metric; import org.apache.kafka.common.MetricName; import org.apache.kafka.common.PartitionInfo; @@ -54,12 +53,6 @@ public interface Producer extends Closeable { void sendOffsetsToTransaction(Map offsets, String consumerGroupId) throws ProducerFencedException; - - /** - * See {@link KafkaProducer#sendOffsetsToTransaction(Map)} - */ - void sendOffsetsToTransaction(Map offsets) throws ProducerFencedException; - /** * See {@link KafkaProducer#commitTransaction()} */ diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java index 14919d054fb30..713e9e4a7662b 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/internals/TransactionManager.java @@ -19,7 +19,6 @@ import org.apache.kafka.clients.ClientResponse; import org.apache.kafka.clients.RequestCompletionHandler; import org.apache.kafka.clients.consumer.OffsetAndMetadata; -import org.apache.kafka.clients.consumer.internals.ConsumerGroupMetadata; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.Node; import org.apache.kafka.common.TopicPartition; @@ -33,7 +32,6 @@ import org.apache.kafka.common.errors.UnsupportedVersionException; import org.apache.kafka.common.message.FindCoordinatorRequestData; import org.apache.kafka.common.message.InitProducerIdRequestData; -import org.apache.kafka.common.message.TxnOffsetCommitRequestData; import org.apache.kafka.common.protocol.Errors; import org.apache.kafka.common.record.DefaultRecordBatch; import org.apache.kafka.common.record.RecordBatch; @@ -202,7 +200,6 @@ public void resetSequenceNumbers(Consumer resetSequence) { private volatile RuntimeException lastError = null; private volatile ProducerIdAndEpoch producerIdAndEpoch; private volatile boolean transactionStarted = false; - private org.apache.kafka.clients.consumer.Consumer consumer; private enum State { UNINITIALIZED, @@ -279,10 +276,6 @@ public void setTransactionalId(String transactionalId) { this.transactionalId = transactionalId; } - public void setConsumer(org.apache.kafka.clients.consumer.Consumer consumer) { - this.consumer = consumer; - } - public synchronized TransactionalRequestResult initializeTransactions() { return handleCachedTransactionRequestResult(() -> { transitionTo(State.INITIALIZING); @@ -998,19 +991,8 @@ private TxnOffsetCommitHandler txnOffsetCommitHandler(TransactionalRequestResult pendingTxnOffsetCommits.put(entry.getKey(), committedOffset); } - - ConsumerGroupMetadata metadata = (ConsumerGroupMetadata) consumer.groupMetadata(); - TxnOffsetCommitRequest.Builder builder = new TxnOffsetCommitRequest.Builder( - new TxnOffsetCommitRequestData() - .setTransactionalId(transactionalId) - .setGroupId(consumerGroupId) - .setProducerId(producerIdAndEpoch.producerId) - .setProducerEpoch(producerIdAndEpoch.epoch) - .setTopics(TxnOffsetCommitRequest.getTopics(pendingTxnOffsetCommits)) - .setGenerationId(metadata.generation()) - .setMemberId(metadata.memberId()) - .setGroupInstanceId(metadata.groupInstanceId().orElse(null)) - ); + TxnOffsetCommitRequest.Builder builder = new TxnOffsetCommitRequest.Builder(transactionalId, consumerGroupId, + producerIdAndEpoch.producerId, producerIdAndEpoch.epoch, pendingTxnOffsetCommits); return new TxnOffsetCommitHandler(result, builder); } @@ -1217,16 +1199,6 @@ public void handleResponse(AbstractResponse response) { } else if (error == Errors.INVALID_PRODUCER_EPOCH) { fatalError(error.exception()); return; - } else if (error == Errors.ILLEGAL_GENERATION) { - fatalError(error.exception()); - return; - } else if (error == Errors.FENCED_INSTANCE_ID) { - fatalError(error.exception()); - return; - } else if (error == Errors.UNKNOWN_MEMBER_ID) { - // What should we do? The consumer internal is not available here. - reenqueue(); - return; } else if (error == Errors.TRANSACTIONAL_ID_AUTHORIZATION_FAILED) { fatalError(error.exception()); return; @@ -1465,7 +1437,7 @@ FindCoordinatorRequest.CoordinatorType coordinatorType() { @Override String coordinatorKey() { - return builder.data.groupId(); + return builder.consumerGroupId(); } @Override @@ -1474,7 +1446,7 @@ public void handleResponse(AbstractResponse response) { boolean coordinatorReloaded = false; Map errors = txnOffsetCommitResponse.errors(); - log.debug("Received TxnOffsetCommit response for consumer group {}: {}", builder.data.groupId(), + log.debug("Received TxnOffsetCommit response for consumer group {}: {}", builder.consumerGroupId(), errors); for (Map.Entry entry : errors.entrySet()) { @@ -1487,14 +1459,14 @@ public void handleResponse(AbstractResponse response) { || error == Errors.REQUEST_TIMED_OUT) { if (!coordinatorReloaded) { coordinatorReloaded = true; - lookupCoordinator(FindCoordinatorRequest.CoordinatorType.GROUP, builder.data.groupId()); + lookupCoordinator(FindCoordinatorRequest.CoordinatorType.GROUP, builder.consumerGroupId()); } } else if (error == Errors.UNKNOWN_TOPIC_OR_PARTITION || error == Errors.COORDINATOR_LOAD_IN_PROGRESS) { // If the topic is unknown or the coordinator is loading, retry with the current coordinator continue; } else if (error == Errors.GROUP_AUTHORIZATION_FAILED) { - abortableError(GroupAuthorizationException.forGroupId(builder.data.groupId())); + abortableError(GroupAuthorizationException.forGroupId(builder.consumerGroupId())); break; } else if (error == Errors.TRANSACTIONAL_ID_AUTHORIZATION_FAILED || error == Errors.INVALID_PRODUCER_EPOCH diff --git a/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java b/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java index a885486062cf3..cca5771002cd8 100644 --- a/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java +++ b/clients/src/test/java/org/apache/kafka/clients/producer/internals/TransactionManagerTest.java @@ -919,7 +919,7 @@ public void testUnsupportedForMessageFormatInTxnOffsetCommit() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.NONE, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled, TxnOffsetCommit Enqueued @@ -1089,7 +1089,7 @@ public void testGroupAuthorizationFailureInFindCoordinator() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(new TopicPartition("foo", 0), new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(new TopicPartition("foo", 0), new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.NONE, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled, TxnOffsetCommit Enqueued @@ -1121,7 +1121,7 @@ public void testGroupAuthorizationFailureInTxnOffsetCommit() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(tp1, new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(tp1, new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.NONE, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled, TxnOffsetCommit Enqueued @@ -1157,7 +1157,7 @@ public void testTransactionalIdAuthorizationFailureInAddOffsetsToTxn() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.TRANSACTIONAL_ID_AUTHORIZATION_FAILED, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled @@ -1182,7 +1182,7 @@ public void testTransactionalIdAuthorizationFailureInTxnOffsetCommit() { transactionManager.beginTransaction(); TransactionalRequestResult sendOffsetsResult = transactionManager.sendOffsetsToTransaction( - singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); + singletonMap(tp, new OffsetAndMetadata(39L)), consumerGroupId); prepareAddOffsetsToTxnResponse(Errors.NONE, consumerGroupId, pid, epoch); sender.runOnce(); // AddOffsetsToTxn Handled, TxnOffsetCommit Enqueued @@ -2826,9 +2826,9 @@ private void prepareTxnOffsetCommitResponse(final String consumerGroupId, Map txnOffsetCommitResponse) { client.prepareResponse(request -> { TxnOffsetCommitRequest txnOffsetCommitRequest = (TxnOffsetCommitRequest) request; - assertEquals(consumerGroupId, txnOffsetCommitRequest.data.groupId()); - assertEquals(producerId, txnOffsetCommitRequest.data.producerId()); - assertEquals(producerEpoch, txnOffsetCommitRequest.data.producerEpoch()); + assertEquals(consumerGroupId, txnOffsetCommitRequest.consumerGroupId()); + assertEquals(producerId, txnOffsetCommitRequest.producerId()); + assertEquals(producerEpoch, txnOffsetCommitRequest.producerEpoch()); return true; }, new TxnOffsetCommitResponse(0, txnOffsetCommitResponse)); } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStandbyTasks.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStandbyTasks.java index 30e4914fc98a8..1d607b21be9a1 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStandbyTasks.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStandbyTasks.java @@ -21,7 +21,7 @@ class AssignedStandbyTasks extends AssignedTasks { AssignedStandbyTasks(final LogContext logContext) { - super(logContext, "standby task", null, null); + super(logContext, "standby task", null, null, "dummy-group-id"); } } diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasks.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasks.java index cfd37d1af45d2..4d44f97f0fcdd 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasks.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasks.java @@ -41,8 +41,9 @@ class AssignedStreamsTasks extends AssignedTasks implements Restorin AssignedStreamsTasks(final LogContext logContext, final Producer threadProducer, - final Time time) { - super(logContext, "stream task", threadProducer, time); + final Time time, + final String consumerGroupId) { + super(logContext, "stream task", threadProducer, time, consumerGroupId); this.threadProducer = threadProducer; } @@ -142,9 +143,8 @@ void addToRestoring(final StreamTask task) { * or if the task producer got fenced (EOS) */ int maybeCommitPerUserRequested() { - return commitInternal(log, - threadProducer, - task -> task.commitRequested() && task.commitNeeded()); + return commitInternal(log, threadProducer, + task -> task.commitRequested() && task.commitNeeded()); } /** diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java index 132520f9c1b99..af7cbdec54669 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java @@ -47,6 +47,7 @@ abstract class AssignedTasks { private final Set previousActiveTasks = new HashSet<>(); private final Producer eosProducer; private final Time time; + private final String consumerGroupId; // IQ may access this map. final Map running = new ConcurrentHashMap<>(); @@ -55,11 +56,13 @@ abstract class AssignedTasks { AssignedTasks(final LogContext logContext, final String taskTypeName, final Producer eosProducer, - final Time time) { + final Time time, + final String consumerGroupId) { this.taskTypeName = taskTypeName; this.log = logContext.logger(getClass()); this.eosProducer = eosProducer; this.time = time; + this.consumerGroupId = consumerGroupId; } void addNewTask(final T task) { @@ -296,14 +299,14 @@ public interface TaskStatus { boolean needsCommit(Task task); } - protected int commitInternal(Logger log, - Producer eosProducer, - TaskStatus taskStatus) { + protected int commitInternal(final Logger log, + final Producer eosProducer, + final TaskStatus taskStatus) { int committed = 0; RuntimeException firstException = null; - Map pendingOffsets = new HashMap<>(); - List tasks = new ArrayList<>(); + final Map pendingOffsets = new HashMap<>(); + final List tasks = new ArrayList<>(); for (final Iterator it = running().iterator(); it.hasNext(); ) { final T task = it.next(); @@ -341,10 +344,10 @@ protected int commitInternal(Logger log, if (!pendingOffsets.isEmpty()) { final long startNs = time.nanoseconds(); - eosProducer.sendOffsetsToTransaction(pendingOffsets); + eosProducer.sendOffsetsToTransaction(pendingOffsets, consumerGroupId); eosProducer.commitTransaction(); - long commitLatency = time.nanoseconds() - startNs; - for (Task task : tasks) { + final long commitLatency = time.nanoseconds() - startNs; + for (final Task task : tasks) { task.markCommitDone(commitLatency); } eosProducer.beginTransaction(); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index 562500f9d4e75..1bbd5c71bc4a4 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -153,7 +153,7 @@ public StreamTask(final TaskId id, final ThreadCache cache, final Time time, final ProducerSupplier producerSupplier, - boolean isTaskProducer) { + final boolean isTaskProducer) { this(id, partitions, topology, consumer, changelogReader, config, metrics, stateDirectory, cache, time, producerSupplier, null, isTaskProducer); } @@ -169,7 +169,7 @@ public StreamTask(final TaskId id, final Time time, final ProducerSupplier producerSupplier, final RecordCollector recordCollector, - boolean isTaskProducer) { + final boolean isTaskProducer) { super(id, partitions, topology, consumer, changelogReader, false, stateDirectory, config); this.isTaskProducer = isTaskProducer; @@ -493,7 +493,7 @@ public Map getPendingOffsets() { } @Override - public void markCommitDone(long commitLatency) { + public void markCommitDone(final long commitLatency) { commitNeeded = false; commitRequested = false; taskMetrics.taskCommitTimeSensor.record(commitLatency); @@ -875,7 +875,7 @@ Producer getProducer() { private void initializeTransactions() { try { - producer.initTransactions(consumer); + producer.initTransactions(); producer.beginTransaction(); } catch (final TimeoutException retriable) { log.error( diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index 1a6cdc1a9ed95..a7c8bddbc2f53 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -393,7 +393,6 @@ Collection createTasks(final Consumer consumer, log.trace("Created task {} with assigned partitions {}", taskId, partitions); createdTasks.add(task); } - } return createdTasks; } @@ -627,6 +626,8 @@ public static StreamThread create(final InternalTopologyBuilder builder, changelogReader, time, log); + + final String applicationId = config.getString(StreamsConfig.APPLICATION_ID_CONFIG); final TaskManager taskManager = new TaskManager( changelogReader, processId, @@ -636,11 +637,10 @@ public static StreamThread create(final InternalTopologyBuilder builder, activeTaskCreator, standbyTaskCreator, adminClient, - new AssignedStreamsTasks(logContext, threadProducer, time), + new AssignedStreamsTasks(logContext, threadProducer, time, applicationId), new AssignedStandbyTasks(logContext)); log.info("Creating consumer client"); - final String applicationId = config.getString(StreamsConfig.APPLICATION_ID_CONFIG); final Map consumerConfigs = config.getMainConsumerConfigs(applicationId, getConsumerClientId(threadClientId), threadIdx); consumerConfigs.put(StreamsConfig.InternalConfig.TASK_MANAGER_FOR_PARTITION_ASSIGNOR, taskManager); final AtomicInteger assignmentErrorCode = new AtomicInteger(); @@ -651,15 +651,15 @@ public static StreamThread create(final InternalTopologyBuilder builder, consumerConfigs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "none"); } - final Consumer mainConsumer = clientSupplier.getConsumer(consumerConfigs); - taskManager.setConsumer(mainConsumer); + final Consumer consumer = clientSupplier.getConsumer(consumerConfigs); + taskManager.setConsumer(consumer); return new StreamThread( time, config, threadProducer, restoreConsumer, - mainConsumer, + consumer, originalReset, taskManager, streamsMetrics, @@ -703,7 +703,7 @@ public StreamThread(final Time time, final String threadClientId, final LogContext logContext, final AtomicInteger assignmentErrorCode, - boolean eosEnabled) { + final boolean eosEnabled) { super(threadClientId); this.stateLock = new Object(); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java index af09fed9c4563..3295ab05ec77e 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/Task.java @@ -22,9 +22,9 @@ import org.apache.kafka.streams.processor.ProcessorContext; import org.apache.kafka.streams.processor.StateStore; import org.apache.kafka.streams.processor.TaskId; -import sun.reflect.generics.reflectiveObjects.NotImplementedException; import java.util.Collection; +import java.util.Collections; import java.util.Map; import java.util.Set; @@ -48,11 +48,11 @@ default boolean commitRequested() { void commit(); default Map getPendingOffsets() { - throw new NotImplementedException(); + return Collections.emptyMap(); } default void markCommitDone(long commitLatency) { - throw new NotImplementedException(); + // No-op } void suspend(); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasksTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasksTest.java index ffd0f8baae6b9..f3594d8c0b484 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasksTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasksTest.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.utils.LogContext; +import org.apache.kafka.common.utils.MockTime; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.streams.errors.TaskMigratedException; import org.apache.kafka.streams.processor.TaskId; @@ -52,7 +53,7 @@ public class AssignedStreamsTasksTest { @Before public void before() { - assignedTasks = new AssignedStreamsTasks(new LogContext("log ")); + assignedTasks = new AssignedStreamsTasks(new LogContext("log"), null, new MockTime(), "dummy-group-id"); EasyMock.expect(t1.id()).andReturn(taskId1).anyTimes(); EasyMock.expect(t2.id()).andReturn(taskId2).anyTimes(); } diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java index 5e768b1000b45..ffe876a5bac68 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java @@ -751,7 +751,7 @@ private void setStreamThread(final StreamThread streamThread) { null, null, null, - new AssignedStreamsTasks(new LogContext()), + new AssignedStreamsTasks(new LogContext(), null, new MockTime(), applicationId), new AssignedStandbyTasks(new LogContext())); taskManager.setConsumer(mockStreamThreadConsumer); taskManager.setAssignmentMetadata(Collections.emptyMap(), Collections.emptyMap()); From 92a8d9302d82f9e8a84b64be25120bd3ba0ca802 Mon Sep 17 00:00:00 2001 From: abbccdda Date: Tue, 6 Aug 2019 15:03:52 -0700 Subject: [PATCH 3/6] simply the initialization logic --- .../apache/kafka/streams/StreamsConfig.java | 40 +++++----------- .../processor/internals/AssignedTasks.java | 6 +-- .../processor/internals/StreamTask.java | 1 - .../processor/internals/StreamThread.java | 38 +++++---------- .../internals/AssignedStreamsTasksTest.java | 8 ++++ .../processor/internals/StreamThreadTest.java | 47 +++++++++++++++++++ 6 files changed, 83 insertions(+), 57 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java index 66a4f4132be2c..455f8722256fd 100644 --- a/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java +++ b/streams/src/main/java/org/apache/kafka/streams/StreamsConfig.java @@ -251,34 +251,6 @@ public class StreamsConfig extends AbstractConfig { @SuppressWarnings("WeakerAccess") public static final String UPGRADE_FROM_11 = "1.1"; - /** - * Config value for parameter {@link #UPGRADE_FROM_CONFIG "upgrade.from"} for upgrading an application from version {@code 2.0.x}. - */ - @SuppressWarnings("WeakerAccess") - public static final String UPGRADE_FROM_20 = "2.0"; - - - /** - * Config value for parameter {@link #UPGRADE_FROM_CONFIG "upgrade.from"} for upgrading an application from version {@code 2.1.x}. - */ - @SuppressWarnings("WeakerAccess") - public static final String UPGRADE_FROM_21 = "2.1"; - - - - /** - * Config value for parameter {@link #UPGRADE_FROM_CONFIG "upgrade.from"} for upgrading an application from version {@code 2.2.x}. - */ - @SuppressWarnings("WeakerAccess") - public static final String UPGRADE_FROM_22 = "2.2"; - - - /** - * Config value for parameter {@link #UPGRADE_FROM_CONFIG "upgrade.from"} for upgrading an application from version {@code 2.3.x}. - */ - @SuppressWarnings("WeakerAccess") - public static final String UPGRADE_FROM_23 = "2.3"; - /** * Config value for parameter {@link #PROCESSING_GUARANTEE_CONFIG "processing.guarantee"} for at-least-once processing guarantees. */ @@ -291,6 +263,13 @@ public class StreamsConfig extends AbstractConfig { @SuppressWarnings("WeakerAccess") public static final String EXACTLY_ONCE = "exactly_once"; + /** + * Config to use thread level producer + */ + @SuppressWarnings("WeakerAccess") + public static final String USE_EOS_THREAD_PRODUCER_CONFIG = "use.eos.thread.producer"; + private static final String USE_EOS_THREAD_PRODUCER_DOC = "Flag to use thread level producer."; + /** {@code application.id} */ @SuppressWarnings("WeakerAccess") public static final String APPLICATION_ID_CONFIG = "application.id"; @@ -596,6 +575,11 @@ public class StreamsConfig extends AbstractConfig { in(NO_OPTIMIZATION, OPTIMIZE), Importance.MEDIUM, TOPOLOGY_OPTIMIZATION_DOC) + .define(USE_EOS_THREAD_PRODUCER_CONFIG, + Type.BOOLEAN, + false, + Importance.MEDIUM, + USE_EOS_THREAD_PRODUCER_DOC) // LOW diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java index af7cbdec54669..40fa3fd08c9c2 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java @@ -306,7 +306,7 @@ protected int commitInternal(final Logger log, RuntimeException firstException = null; final Map pendingOffsets = new HashMap<>(); - final List tasks = new ArrayList<>(); + final List externalCommitTasks = new ArrayList<>(); for (final Iterator it = running().iterator(); it.hasNext(); ) { final T task = it.next(); @@ -314,7 +314,7 @@ protected int commitInternal(final Logger log, if (taskStatus.needsCommit(task)) { if (eosProducer != null) { pendingOffsets.putAll(task.getPendingOffsets()); - tasks.add(task); + externalCommitTasks.add(task); } else { task.commit(); } @@ -347,7 +347,7 @@ protected int commitInternal(final Logger log, eosProducer.sendOffsetsToTransaction(pendingOffsets, consumerGroupId); eosProducer.commitTransaction(); final long commitLatency = time.nanoseconds() - startNs; - for (final Task task : tasks) { + for (final Task task : externalCommitTasks) { task.markCommitDone(commitLatency); } eosProducer.beginTransaction(); diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index 1bbd5c71bc4a4..4f08970390d5e 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -876,7 +876,6 @@ Producer getProducer() { private void initializeTransactions() { try { producer.initTransactions(); - producer.beginTransaction(); } catch (final TimeoutException retriable) { log.error( "Timeout exception caught when initializing transactions for task {}. " + diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index a7c8bddbc2f53..d899f013314e4 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -592,16 +592,25 @@ public static StreamThread create(final InternalTopologyBuilder builder, final Duration pollTime = Duration.ofMillis(config.getLong(StreamsConfig.POLL_MS_CONFIG)); final StoreChangelogReader changelogReader = new StoreChangelogReader(restoreConsumer, pollTime, userStateRestoreListener, logContext); - Producer threadProducer = null; final boolean eosEnabled = StreamsConfig.EXACTLY_ONCE.equals(config.getString(StreamsConfig.PROCESSING_GUARANTEE_CONFIG)); - final String upgradeFrom = config.getString(StreamsConfig.UPGRADE_FROM_CONFIG); - if (useThreadProducer(eosEnabled, upgradeFrom)) { + final boolean useEosThreadProducer = config.getBoolean(StreamsConfig.USE_EOS_THREAD_PRODUCER_CONFIG); + + final Producer threadProducer; + final String applicationId = config.getString(StreamsConfig.APPLICATION_ID_CONFIG); + if (useEosThreadProducer || !eosEnabled) { final Map producerConfigs = config.getProducerConfigs(getThreadProducerClientId(threadClientId)); log.info("Creating shared producer client"); + if (eosEnabled) { + producerConfigs.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, applicationId + "-" + threadClientId); + } threadProducer = clientSupplier.getProducer(producerConfigs); + } else { + threadProducer = null; } + final Producer eosThreadProducer = eosEnabled ? threadProducer : null; + final StreamsMetricsImpl streamsMetrics = new StreamsMetricsImpl(metrics, threadClientId); final ThreadCache cache = new ThreadCache(logContext, cacheSizeBytes, streamsMetrics); @@ -627,7 +636,6 @@ public static StreamThread create(final InternalTopologyBuilder builder, time, log); - final String applicationId = config.getString(StreamsConfig.APPLICATION_ID_CONFIG); final TaskManager taskManager = new TaskManager( changelogReader, processId, @@ -637,7 +645,7 @@ public static StreamThread create(final InternalTopologyBuilder builder, activeTaskCreator, standbyTaskCreator, adminClient, - new AssignedStreamsTasks(logContext, threadProducer, time, applicationId), + new AssignedStreamsTasks(logContext, eosThreadProducer, time, applicationId), new AssignedStandbyTasks(logContext)); log.info("Creating consumer client"); @@ -671,26 +679,6 @@ public static StreamThread create(final InternalTopologyBuilder builder, .updateThreadMetadata(getSharedAdminClientId(clientId)); } - private static boolean useThreadProducer(final boolean eosEnabled, final String upgradeFrom) { - switch (upgradeFrom) { - case StreamsConfig.UPGRADE_FROM_0100: - case StreamsConfig.UPGRADE_FROM_0101: - case StreamsConfig.UPGRADE_FROM_0102: - case StreamsConfig.UPGRADE_FROM_0110: - case StreamsConfig.UPGRADE_FROM_10: - case StreamsConfig.UPGRADE_FROM_11: - case StreamsConfig.UPGRADE_FROM_20: - case StreamsConfig.UPGRADE_FROM_21: - case StreamsConfig.UPGRADE_FROM_22: - case StreamsConfig.UPGRADE_FROM_23: - // If upgrading from an older version, we shall continue using the eos producer setup. - return !eosEnabled; - default: - // If not upgrading, start from 2.4 there will be no task level producer setup. - return true; - } - } - public StreamThread(final Time time, final StreamsConfig config, final Producer producer, diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasksTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasksTest.java index f3594d8c0b484..d07597c35337c 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasksTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasksTest.java @@ -452,6 +452,14 @@ public void shouldReturnNumberOfPunctuations() { EasyMock.verify(t1); } + @Test + public void shouldLetThreadProducerCommit() { +// final Map producerConfig = new HashMap<>(); +// +// assignedTasks = new AssignedStreamsTasks(new LogContext("log"), new DefaultKafkaClientSupplier().getProducer(), new MockTime(), "dummy-group-id"); + + } + private void addAndInitTask() { assignedTasks.addNewTask(t1); assignedTasks.initializeNewTasks(); diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java index ffe876a5bac68..25c42272c2d1a 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamThreadTest.java @@ -106,6 +106,7 @@ import static org.junit.Assert.assertNotEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -452,7 +453,53 @@ public void shouldRespectNumIterationsInMainLoop() { thread.runOnce(); assertThat(thread.currentNumIterations(), equalTo(1)); + } + + @Test + public void threadProducerOnEosShouldCommit() { + final MockProcessor mockProcessor = new MockProcessor(PunctuationType.WALL_CLOCK_TIME, 10L); + internalTopologyBuilder.addSource(null, "source1", null, null, null, topic1); + internalTopologyBuilder.addProcessor("processor1", () -> mockProcessor, "source1"); + internalTopologyBuilder.addProcessor("processor2", () -> new MockProcessor(PunctuationType.STREAM_TIME, 10L), "source1"); + + final Properties properties = new Properties(); + properties.put(StreamsConfig.COMMIT_INTERVAL_MS_CONFIG, 100L); + properties.put(StreamsConfig.USE_EOS_THREAD_PRODUCER_CONFIG, true); + properties.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE); + final StreamsConfig config = new StreamsConfig(StreamsTestUtils.getStreamsConfig(applicationId, + "localhost:2171", + Serdes.ByteArraySerde.class.getName(), + Serdes.ByteArraySerde.class.getName(), + properties)); + final StreamThread thread = createStreamThread(clientId, config, true); + + thread.setState(StreamThread.State.STARTING); + thread.setState(StreamThread.State.PARTITIONS_REVOKED); + + final Set assignedPartitions = Collections.singleton(t1p1); + thread.taskManager().setAssignmentMetadata( + Collections.singletonMap( + new TaskId(0, t1p1.partition()), + assignedPartitions), + Collections.emptyMap()); + + final MockConsumer mockConsumer = (MockConsumer) thread.consumer; + mockConsumer.assign(Collections.singleton(t1p1)); + mockConsumer.updateBeginningOffsets(Collections.singletonMap(t1p1, 0L)); + thread.rebalanceListener.onPartitionsAssigned(assignedPartitions); + thread.runOnce(); + // processed one record, punctuated after the first record, and hence num.iterations is still 1 + long offset = -1; + addRecord(mockConsumer, ++offset, 0L); + thread.runOnce(); + + assertThat(thread.currentNumIterations(), equalTo(1)); + + mockProcessor.requestCommit(); + addRecord(mockConsumer, ++offset, 15L); + // Mock producer should throw because it's not enabled for transactional support but wanted. + assertThrows(IllegalStateException.class, () -> thread.runOnce()); } @Test From 41f80dddcd7b9cadf99a8261a44c61e8298d99bd Mon Sep 17 00:00:00 2001 From: abbccdda Date: Tue, 6 Aug 2019 22:43:13 -0700 Subject: [PATCH 4/6] fix stream task test --- .../streams/processor/internals/StreamTaskTest.java | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java index b87ab08bd2354..836c231ee479d 100644 --- a/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java +++ b/streams/src/test/java/org/apache/kafka/streams/processor/internals/StreamTaskTest.java @@ -1506,18 +1506,19 @@ private StreamTask createStatefulTaskThatThrowsExceptionOnClose() { singletonList(stateStore), Collections.emptyMap()); + final boolean enableEOS = true; return new StreamTask( taskId00, partitions, topology, consumer, changelogReader, - createConfig(true), + createConfig(enableEOS), streamsMetrics, stateDirectory, null, time, - () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer), false); + () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer), enableEOS); } private StreamTask createStatelessTask(final StreamsConfig streamsConfig) { @@ -1531,6 +1532,8 @@ private StreamTask createStatelessTask(final StreamsConfig streamsConfig) { source1.addChild(processorSystemTime); source2.addChild(processorSystemTime); + final boolean enableEOS = StreamsConfig.EXACTLY_ONCE.equals(streamsConfig.getString(StreamsConfig.PROCESSING_GUARANTEE_CONFIG)); + return new StreamTask( taskId00, partitions, @@ -1542,7 +1545,7 @@ private StreamTask createStatelessTask(final StreamsConfig streamsConfig) { stateDirectory, null, time, - () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer), false); + () -> producer = new MockProducer<>(false, bytesSerializer, bytesSerializer), enableEOS); } // this task will throw exception when processing (on partition2), flushing, suspending and closing From ca78fa7b3e1516fea3f27ff1f07156aa323291c7 Mon Sep 17 00:00:00 2001 From: abbccdda Date: Wed, 7 Aug 2019 16:18:09 -0700 Subject: [PATCH 5/6] add debug log and transaction beginning --- .../processor/internals/StreamTask.java | 20 ++++++++++++++----- .../processor/internals/StreamThread.java | 14 ++++++++++++- 2 files changed, 28 insertions(+), 6 deletions(-) diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java index 4f08970390d5e..1ffa77929595d 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamTask.java @@ -183,14 +183,22 @@ public StreamTask(final TaskId id, final ProductionExceptionHandler productionExceptionHandler = config.defaultProductionExceptionHandler(); if (recordCollector == null) { + log.info("record collector is initialized on task"); this.recordCollector = new RecordCollectorImpl( id.toString(), logContext, productionExceptionHandler, ThreadMetrics.skipRecordSensor(streamsMetrics)); } else { + log.info("record collector given is non-null"); this.recordCollector = recordCollector; } + + if (this.producer == null) { + log.info("Initializing the record collector with producer null"); + } else { + log.info("Initializing the record collector with non-null producer"); + } this.recordCollector.init(this.producer); streamTimePunctuationQueue = new PunctuationQueue(); @@ -282,12 +290,14 @@ public void initializeTopology() { @Override public void resume() { log.debug("Resuming"); - if (isTaskProducer) { - if (producer != null) { - throw new IllegalStateException("Task producer should be null."); - } + if (eosEnabled) { +// if (producer != null) { +// throw new IllegalStateException("Task producer should be null."); +// } producer = producerSupplier.get(); - initializeTransactions(); + if (isTaskProducer) { + initializeTransactions(); + } recordCollector.init(producer); try { diff --git a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java index d899f013314e4..bf219f4df94fb 100644 --- a/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java +++ b/streams/src/main/java/org/apache/kafka/streams/processor/internals/StreamThread.java @@ -431,6 +431,11 @@ static class TaskCreator extends AbstractTaskCreator { this.cache = cache; this.clientSupplier = clientSupplier; this.threadProducer = threadProducer; + if (threadProducer == null) { + log.info("Initializing with thread producer to null in task creator"); + } else { + log.info("Initializing with thread producer to non-null in task creator"); + } this.threadClientId = threadClientId; createTaskSensor = ThreadMetrics.createTaskSensor(streamsMetrics); } @@ -440,7 +445,7 @@ StreamTask createTask(final Consumer consumer, final TaskId taskId, final Set partitions) { createTaskSensor.record(); - + log.info("creating task {}", taskId); return new StreamTask( taskId, partitions, @@ -615,6 +620,12 @@ public static StreamThread create(final InternalTopologyBuilder builder, final ThreadCache cache = new ThreadCache(logContext, cacheSizeBytes, streamsMetrics); + if (threadProducer == null) { + log.info("Init thread producer on stream-thread with null"); + } else { + log.info("Init thread producer on stream-thread with non-null"); + } + final AbstractTaskCreator activeTaskCreator = new TaskCreator( builder, config, @@ -800,6 +811,7 @@ private void maybeInitializeTransactions() { // This is a thread-level txn producer if (eosEnabled && producer != null) { producer.initTransactions(); + producer.beginTransaction(); } } catch (final TimeoutException retriable) { log.error( From aa582d02334ea6db0ccafececa7a2711b624f5d0 Mon Sep 17 00:00:00 2001 From: abbccdda Date: Wed, 7 Aug 2019 21:53:04 -0700 Subject: [PATCH 6/6] system test --- tests/kafkatest/tests/streams/base_streams_test.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/tests/kafkatest/tests/streams/base_streams_test.py b/tests/kafkatest/tests/streams/base_streams_test.py index 53e4231edb419..2b2e92fdb6ec1 100644 --- a/tests/kafkatest/tests/streams/base_streams_test.py +++ b/tests/kafkatest/tests/streams/base_streams_test.py @@ -67,7 +67,7 @@ def assert_produce(self, topic, test_state, num_messages=5, timeout_sec=60): wait_until(lambda: producer.num_acked >= num_messages, timeout_sec=timeout_sec, - err_msg="At %s failed to send messages " % test_state) + err_msg="At %s failed to send messages. Expected: %s, actual: %s " % (test_state, num_messages, producer.num_acked)) def assert_consume(self, client_id, test_state, topic, num_messages=5, timeout_sec=60): consumer = self.get_consumer(client_id, topic, num_messages) @@ -75,7 +75,7 @@ def assert_consume(self, client_id, test_state, topic, num_messages=5, timeout_s wait_until(lambda: consumer.total_consumed() >= num_messages, timeout_sec=timeout_sec, - err_msg="At %s streams did not process messages in %s seconds " % (test_state, timeout_sec)) + err_msg="At %s streams did not process messages in %s seconds. Expected: %s, actual: %s " % (test_state, timeout_sec, num_messages, consumer.total_consumed())) @staticmethod def get_configs(extra_configs=""):