This is an automated email from the ASF dual-hosted git repository.
1996fanrui pushed a change to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git
from dc707db6 [FLINK-39888][Kafka] Allow configuring offset reset strategy
independently (#267)
new 564ad828 [FLINK-40128][connector] Introduce topic integrity utilities
in static kafka source
new 9cc364d8 [FLINK-40128][connector] Add an opt-in based topic integrity
config option
new a5b7bc38 [FLINK-40128][connector] Implement topic integrity awareness
across kafka subscribers
new db222b3c [FLINK-40128][connector] Topic integrity ITCase
The 4 revisions listed above as "new" are entirely new to this
repository and will be described in separate emails. The revisions
listed as "add" were already present in the repository and have only
been added to this reference.
Summary of changes:
.../content.zh/docs/connectors/datastream/kafka.md | 32 +++
docs/content/docs/connectors/datastream/kafka.md | 37 +++
.../c0d94764-76a0-4c50-b617-70b1754c4612 | 18 +-
flink-connector-kafka/pom.xml | 8 +
.../connector/kafka/source/KafkaSourceBuilder.java | 26 ++
.../connector/kafka/source/KafkaSourceOptions.java | 7 +
.../source/enumerator/KafkaSourceEnumState.java | 25 ++
.../enumerator/KafkaSourceEnumStateSerializer.java | 75 +++++-
.../source/enumerator/KafkaSourceEnumerator.java | 38 ++-
.../metadata/TopicIntegrityException.java} | 22 +-
.../metadata/TopicIntegrityProvider.java | 170 +++++++++++++
.../enumerator/metadata/TopicMetadataProvider.java | 58 +++++
.../metadata/TopicMetadataSettable.java} | 10 +-
.../subscriber/PartitionSetSubscriber.java | 23 +-
.../enumerator/subscriber/TopicListSubscriber.java | 23 +-
.../subscriber/TopicPatternSubscriber.java | 23 +-
.../kafka/source/KafkaSourceBuilderTest.java | 51 ++++
.../kafka/source/SourceTopicIntegrityTest.java | 280 +++++++++++++++++++++
.../KafkaSourceEnumStateSerializerTest.java | 42 +++-
.../enumerator/KafkaSourceEnumeratorTest.java | 41 +++
.../enumerator/TopicIntegrityProviderTest.java | 195 ++++++++++++++
.../enumerator/subscriber/KafkaSubscriberTest.java | 76 +++++-
.../kafka/testutils/KafkaSourceTestEnv.java | 71 +++++-
23 files changed, 1297 insertions(+), 54 deletions(-)
copy
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/{sink/TopicSelector.java
=> source/enumerator/metadata/TopicIntegrityException.java} (67%)
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/metadata/TopicIntegrityProvider.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/metadata/TopicMetadataProvider.java
copy
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/{dynamic/source/GetMetadataUpdateEvent.java
=> source/enumerator/metadata/TopicMetadataSettable.java} (75%)
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/SourceTopicIntegrityTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/TopicIntegrityProviderTest.java