This is an automated email from the ASF dual-hosted git repository.
martijnvisser pushed a change to branch
dependabot/maven/org.apache.commons-commons-compress-1.24.0
in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git
discard f701ac25 Bump org.apache.commons:commons-compress from 1.22 to 1.24.0
add 825052f5 [FLINK-33361][connectors/kafka] Add Java 17 compatibility to
Flink Kafka connector
add eaeb7817 [FLINK-32416] initial implementation of DynamicKafkaSource
with bounded/unbounded support and unit/integration tests
add cdfa328b [FLINK-32416] Fix flaky tests by ensuring test utilities
produce records with consistency and cleanup notify no more splits to ensure it
is sent. This closes #79
add d3bda90d [hotfix] Synchronize CI pipeline setup. This closes #78
add 8ab7d004 Bump org.apache.commons:commons-compress from 1.22 to 1.24.0
add f3ca94fb [FLINK-33329] Bump org.apache.commons:commons-compress from
1.22 to 1.24.0
This update added new revisions after undoing existing revisions.
That is to say, some revisions that were in the old version of the
branch are not in the new version. This situation occurs
when a user --force pushes a change and generates a repository
containing something like this:
* -- * -- B -- O -- O -- O (f701ac25)
\
N -- N -- N
refs/heads/dependabot/maven/org.apache.commons-commons-compress-1.24.0
(f3ca94fb)
You should already have received notification emails for all of the O
revisions, and so the following emails describe only the N revisions
from the common base, B.
Any revisions marked "omit" are not gone; other references still
refer to them. Any revisions marked "discard" are gone forever.
No new revisions were added by this update.
Summary of changes:
.github/workflows/push_pr.yml | 12 +-
.github/workflows/weekly.yml | 20 +-
.../docs/connectors/table/dynamic-kafka.md | 141 +++
.../docs/connectors/datastream/dynamic-kafka.md | 141 +++
.../flink-end-to-end-tests-common-kafka/pom.xml | 4 +-
.../984f05c0-ec82-405e-9bcc-d202dbe7202e | 357 ++++++++
flink-connector-kafka/pom.xml | 13 +
.../kafka/dynamic/metadata/ClusterMetadata.java | 92 ++
.../dynamic/metadata/KafkaMetadataService.java | 53 ++
.../kafka/dynamic/metadata/KafkaStream.java | 94 ++
.../SingleClusterTopicMetadataService.java | 118 +++
.../kafka/dynamic/source/DynamicKafkaSource.java | 222 +++++
.../dynamic/source/DynamicKafkaSourceBuilder.java | 328 +++++++
.../dynamic/source/DynamicKafkaSourceOptions.java | 60 ++
.../source/GetMetadataUpdateEvent.java} | 18 +-
.../kafka/dynamic/source/MetadataUpdateEvent.java | 77 ++
.../enumerator/DynamicKafkaSourceEnumState.java | 58 ++
.../DynamicKafkaSourceEnumStateSerializer.java | 187 ++++
.../enumerator/DynamicKafkaSourceEnumerator.java | 546 ++++++++++++
.../enumerator/StoppableKafkaEnumContextProxy.java | 316 +++++++
.../subscriber/KafkaStreamSetSubscriber.java} | 28 +-
.../subscriber/KafkaStreamSubscriber.java} | 25 +-
.../subscriber/StreamPatternSubscriber.java | 53 ++
.../source/metrics/KafkaClusterMetricGroup.java | 142 +++
.../metrics/KafkaClusterMetricGroupManager.java | 76 ++
.../source/reader/DynamicKafkaSourceReader.java | 549 ++++++++++++
.../reader/KafkaPartitionSplitReaderWrapper.java | 98 +++
.../source/split/DynamicKafkaSourceSplit.java | 85 ++
.../split/DynamicKafkaSourceSplitSerializer.java} | 48 +-
.../kafka/source/KafkaPropertiesUtil.java | 67 ++
.../source/enumerator/KafkaSourceEnumState.java | 2 +-
.../dynamic/source/DynamicKafkaSourceITTest.java | 694 +++++++++++++++
.../DynamicKafkaSourceEnumStateSerializerTest.java | 118 +++
.../DynamicKafkaSourceEnumeratorTest.java | 966 +++++++++++++++++++++
.../StoppableKafkaEnumContextProxyTest.java | 211 +++++
.../SingleClusterTopicMetadataServiceTest.java | 117 +++
.../metrics/KafkaClusterMetricGroupTest.java | 95 ++
.../reader/DynamicKafkaSourceReaderTest.java | 347 ++++++++
.../DynamicKafkaSourceSplitSerializerTest.java | 47 +
.../DynamicKafkaSourceExternalContext.java | 263 ++++++
... DynamicKafkaSourceExternalContextFactory.java} | 40 +-
.../kafka/testutils/MockKafkaMetadataService.java | 93 ++
.../kafka/testutils/TwoKafkaContainers.java | 62 ++
.../kafka/testutils/YamlFileMetadataService.java | 361 ++++++++
.../testutils/YamlFileMetadataServiceTest.java | 79 ++
.../kafka/DynamicKafkaSourceTestHelper.java | 231 +++++
.../streaming/connectors/kafka/KafkaTestBase.java | 89 +-
.../connectors/kafka/KafkaTestEnvironment.java | 2 +
.../connectors/kafka/KafkaTestEnvironmentImpl.java | 7 +-
.../src/test/resources/log4j2-test.properties | 3 +
.../src/test/resources/stream-metadata.yaml | 19 +
pom.xml | 16 +-
tools/maven/checkstyle.xml | 4 -
53 files changed, 7791 insertions(+), 103 deletions(-)
create mode 100644 docs/content.zh/docs/connectors/table/dynamic-kafka.md
create mode 100644 docs/content/docs/connectors/datastream/dynamic-kafka.md
create mode 100644
flink-connector-kafka/archunit-violations/984f05c0-ec82-405e-9bcc-d202dbe7202e
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/metadata/ClusterMetadata.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/metadata/KafkaMetadataService.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/metadata/KafkaStream.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/metadata/SingleClusterTopicMetadataService.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSource.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceBuilder.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceOptions.java
copy
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/{sink/HeaderProvider.java
=> dynamic/source/GetMetadataUpdateEvent.java} (66%)
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/MetadataUpdateEvent.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumState.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumStateSerializer.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumerator.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/StoppableKafkaEnumContextProxy.java
copy
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/{source/enumerator/TopicPartitionAndAssignmentStatus.java
=> dynamic/source/enumerator/subscriber/KafkaStreamSetSubscriber.java} (54%)
copy
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/{source/enumerator/initializer/OffsetsInitializerValidator.java
=> dynamic/source/enumerator/subscriber/KafkaStreamSubscriber.java} (51%)
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/subscriber/StreamPatternSubscriber.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/metrics/KafkaClusterMetricGroup.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/metrics/KafkaClusterMetricGroupManager.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReader.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/KafkaPartitionSplitReaderWrapper.java
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/split/DynamicKafkaSourceSplit.java
copy
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/{source/split/KafkaPartitionSplitSerializer.java
=> dynamic/source/split/DynamicKafkaSourceSplitSerializer.java} (50%)
create mode 100644
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaPropertiesUtil.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceITTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumStateSerializerTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumeratorTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/StoppableKafkaEnumContextProxyTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/metadata/SingleClusterTopicMetadataServiceTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/metrics/KafkaClusterMetricGroupTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReaderTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/split/DynamicKafkaSourceSplitSerializerTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/DynamicKafkaSourceExternalContext.java
copy
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/{KafkaSourceExternalContextFactory.java
=> DynamicKafkaSourceExternalContextFactory.java} (57%)
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/MockKafkaMetadataService.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/TwoKafkaContainers.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/YamlFileMetadataService.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/YamlFileMetadataServiceTest.java
create mode 100644
flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/DynamicKafkaSourceTestHelper.java
create mode 100644
flink-connector-kafka/src/test/resources/stream-metadata.yaml