You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@kafka.apache.org by gu...@apache.org on 2020/02/05 05:07:27 UTC
[kafka] branch trunk updated (a16dfe6 -> 4090f9a)
This is an automated email from the ASF dual-hosted git repository.
guozhang pushed a change to branch trunk
in repository https://gitbox.apache.org/repos/asf/kafka.git.
from a16dfe6 MINOR: fix checkstyle issue in ConsumerConfig.java (#8038)
add 4090f9a KAFKA-9113: Clean up task management and state management (#7997)
No new revisions were added by this update.
Summary of changes:
checkstyle/suppressions.xml | 2 +
.../consumer/internals/AbstractCoordinator.java | 2 +-
.../consumer/internals/ConsumerCoordinator.java | 28 +-
.../apache/kafka/clients/producer/Callback.java | 1 +
.../kafka/clients/producer/KafkaProducer.java | 8 +-
.../kafka/clients/producer/MockProducer.java | 13 +-
.../org/apache/kafka/streams/KafkaStreams.java | 31 +-
.../org/apache/kafka/streams/StreamsConfig.java | 2 +
.../streams/errors/TaskMigratedException.java | 58 +-
.../internals/AbstractProcessorContext.java | 2 +-
.../streams/processor/internals/AbstractTask.java | 223 +--
.../processor/internals/AssignedStandbyTasks.java | 90 --
.../processor/internals/AssignedStreamsTasks.java | 557 -------
.../streams/processor/internals/AssignedTasks.java | 297 ----
.../processor/internals/ChangelogReader.java | 37 +-
...{Checkpointable.java => ChangelogRegister.java} | 17 +-
.../internals/GlobalStateManagerImpl.java | 39 +-
.../processor/internals/GlobalStateUpdateTask.java | 4 +-
.../internals/InternalTopologyBuilder.java | 39 +-
.../processor/internals/PartitionGroup.java | 2 +-
.../internals/ProcessorNodePunctuator.java | 2 +-
.../processor/internals/ProcessorStateManager.java | 613 +++----
.../processor/internals/RecordCollector.java | 9 +-
.../processor/internals/RecordCollectorImpl.java | 373 +++--
.../processor/internals/RecordDeserializer.java | 4 +-
.../streams/processor/internals/RecordQueue.java | 3 +-
.../internals/RecoverableClientException.java | 33 -
.../streams/processor/internals/StandbyTask.java | 285 ++--
.../streams/processor/internals/StateManager.java | 19 +-
.../processor/internals/StateManagerUtil.java | 144 +-
.../streams/processor/internals/StateRestorer.java | 139 --
.../processor/internals/StoreChangelogReader.java | 932 ++++++++---
.../streams/processor/internals/StreamTask.java | 887 +++++------
.../streams/processor/internals/StreamThread.java | 351 ++--
.../internals/StreamsPartitionAssignor.java | 122 +-
.../internals/StreamsRebalanceListener.java | 105 +-
.../kafka/streams/processor/internals/Task.java | 166 +-
.../streams/processor/internals/TaskManager.java | 768 ++++-----
.../assignment/AssignorConfiguration.java | 72 +-
.../streams/state/internals/RocksDBStore.java | 1 +
.../internals/StreamThreadStateStoreProvider.java | 20 +-
.../org/apache/kafka/streams/KafkaStreamsTest.java | 6 +-
.../streams/integration/EosIntegrationTest.java | 15 +-
.../integration/GlobalThreadShutDownOrderTest.java | 13 +-
.../integration/LagFetchIntegrationTest.java | 3 +-
.../OptimizedKTableIntegrationTest.java | 103 +-
.../ResetPartitionTimeIntegrationTest.java | 2 +-
.../integration/RestoreIntegrationTest.java | 10 +-
.../SmokeTestDriverIntegrationTest.java | 2 +-
.../integration/StoreQueryIntegrationTest.java | 35 +-
.../StreamStreamJoinIntegrationTest.java | 18 +-
.../StreamsUpgradeTestIntegrationTest.java | 1 +
.../integration/TableTableJoinIntegrationTest.java | 4 +-
.../SessionWindowedCogroupedKStreamImplTest.java | 2 -
.../processor/internals/AbstractTaskTest.java | 242 ---
.../internals/AssignedStreamsTasksTest.java | 705 --------
.../internals/GlobalStateManagerImplTest.java | 142 +-
.../processor/internals/GlobalStateTaskTest.java | 4 +-
.../internals/InternalTopologyBuilderTest.java | 4 +-
.../processor/internals/MockChangelogReader.java | 49 +-
.../internals/ProcessorContextImplTest.java | 1 -
.../internals/ProcessorStateManagerTest.java | 1012 +++++-------
.../processor/internals/PunctuationQueueTest.java | 12 +-
.../processor/internals/RecordCollectorTest.java | 988 ++++++++----
.../processor/internals/RecordQueueTest.java | 13 +-
.../streams/processor/internals/SinkNodeTest.java | 58 +-
.../processor/internals/StandbyTaskTest.java | 826 +++-------
.../processor/internals/StateConsumerTest.java | 8 +-
.../processor/internals/StateManagerStub.java | 18 +-
.../processor/internals/StateRestorerTest.java | 111 --
.../internals/StoreChangelogReaderTest.java | 1461 ++++++++---------
.../processor/internals/StreamTaskTest.java | 1679 +++++++-------------
.../processor/internals/StreamThreadTest.java | 695 ++++----
.../internals/StreamsPartitionAssignorTest.java | 62 +-
.../processor/internals/TaskManagerTest.java | 1033 +++++++-----
.../streams/processor/internals/TaskSuite.java | 17 +-
.../streams/state/KeyValueStoreTestDriver.java | 17 +-
.../StreamThreadStateStoreProviderTest.java | 50 +-
.../kafka/streams/tests/StreamsUpgradeTest.java | 28 +-
.../apache/kafka/test/GlobalStateManagerStub.java | 18 +-
.../test/MockBatchingStateRestoreListener.java | 4 -
.../org/apache/kafka/test/MockKeyValueStore.java | 23 +-
.../org/apache/kafka/test/MockRecordCollector.java | 13 +-
.../kafka/test/MockStateRestoreListener.java | 16 +-
.../org/apache/kafka/test/StreamsTestUtils.java | 9 -
.../apache/kafka/streams/TopologyTestDriver.java | 63 +-
.../streams/streams_broker_compatibility_test.py | 8 +-
87 files changed, 6673 insertions(+), 9430 deletions(-)
delete mode 100644 streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStandbyTasks.java
delete mode 100644 streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasks.java
delete mode 100644 streams/src/main/java/org/apache/kafka/streams/processor/internals/AssignedTasks.java
rename streams/src/main/java/org/apache/kafka/streams/processor/internals/{Checkpointable.java => ChangelogRegister.java} (69%)
delete mode 100644 streams/src/main/java/org/apache/kafka/streams/processor/internals/RecoverableClientException.java
delete mode 100644 streams/src/main/java/org/apache/kafka/streams/processor/internals/StateRestorer.java
delete mode 100644 streams/src/test/java/org/apache/kafka/streams/processor/internals/AbstractTaskTest.java
delete mode 100644 streams/src/test/java/org/apache/kafka/streams/processor/internals/AssignedStreamsTasksTest.java
delete mode 100644 streams/src/test/java/org/apache/kafka/streams/processor/internals/StateRestorerTest.java