This is an automated email from the ASF dual-hosted git repository.
bowenli86 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git
The following commit(s) were added to refs/heads/main by this push:
new dc707db6 [FLINK-39888][Kafka] Allow configuring offset reset strategy
independently (#267)
dc707db6 is described below
commit dc707db60fcf3c0637114503c27685c881519c2c
Author: jigar-bhati <[email protected]>
AuthorDate: Mon Aug 17 18:41:09 2026 -0700
[FLINK-39888][Kafka] Allow configuring offset reset strategy independently
(#267)
This PR makes an explicitly configured Kafka auto.offset.reset property
take precedence over the reset strategy inferred from the source's
starting-offset initializer. If the property is not configured, the existing
initializer-derived behavior is preserved.
---
.../content.zh/docs/connectors/datastream/kafka.md | 4 +-
docs/content.zh/docs/connectors/table/kafka.md | 2 +-
.../docs/connectors/table/upsert-kafka.md | 2 +-
docs/content/docs/connectors/datastream/kafka.md | 9 +-
docs/content/docs/connectors/table/kafka.md | 2 +-
docs/content/docs/connectors/table/upsert-kafka.md | 2 +-
.../c0d94764-76a0-4c50-b617-70b1754c4612 | 2 +-
.../kafka/dynamic/source/DynamicKafkaSource.java | 7 +-
.../dynamic/source/DynamicKafkaSourceBuilder.java | 30 +++++-
.../enumerator/DynamicKafkaSourceEnumerator.java | 9 +-
.../source/reader/DynamicKafkaSourceReader.java | 32 ++++--
.../kafka/source/KafkaPropertiesUtil.java | 58 +++++++++++
.../connector/kafka/source/KafkaSourceBuilder.java | 37 +++++--
.../kafka/table/DynamicKafkaTableSource.java | 74 +++++++-------
.../kafka/table/KafkaConnectorOptionsUtil.java | 32 ++++++
.../connectors/kafka/table/KafkaDynamicSource.java | 81 +++++++---------
.../source/DynamicKafkaSourceBuilderTest.java | 102 ++++++++++++++++++++
.../reader/DynamicKafkaSourceReaderTest.java | 10 +-
.../kafka/source/KafkaPropertiesUtilTest.java | 107 +++++++++++++++++++++
.../kafka/source/KafkaSourceBuilderTest.java | 45 +++++++++
.../kafka/table/DynamicKafkaTableFactoryTest.java | 13 +++
.../kafka/table/KafkaDynamicTableFactoryTest.java | 85 +++++++++++++---
22 files changed, 606 insertions(+), 139 deletions(-)
diff --git a/docs/content.zh/docs/connectors/datastream/kafka.md
b/docs/content.zh/docs/connectors/datastream/kafka.md
index c26d819a..26c1f2ac 100644
--- a/docs/content.zh/docs/connectors/datastream/kafka.md
+++ b/docs/content.zh/docs/connectors/datastream/kafka.md
@@ -223,8 +223,8 @@ Kafka Source 支持流式和批式两种运行模式。默认情况下,KafkaSo
Kafka consumer 的配置可以参考 [Apache Kafka
文档](http://kafka.apache.org/documentation/#consumerconfigs)。
-请注意,即使指定了以下配置项,构建器也会将其覆盖:
-- ```auto.offset.reset.strategy``` 被
OffsetsInitializer#getAutoOffsetResetStrategy() 覆盖
+请注意,构建器会设置以下配置项:
+- 如果未显式配置 ```auto.offset.reset```,则会基于
OffsetsInitializer#getAutoOffsetResetStrategy() 设置该配置。起始 offset
初始化器用于选择初始读取位置,而显式配置的 ```auto.offset.reset``` 用于控制已初始化的位置随后变得不可用时的处理方式。
- ```partition.discovery.interval.ms``` 会在批模式下被覆盖为 -1
### 动态分区检查
diff --git a/docs/content.zh/docs/connectors/table/kafka.md
b/docs/content.zh/docs/connectors/table/kafka.md
index 8eba82f1..808fc31d 100644
--- a/docs/content.zh/docs/connectors/table/kafka.md
+++ b/docs/content.zh/docs/connectors/table/kafka.md
@@ -223,7 +223,7 @@ CREATE TABLE KafkaTable (
<td style="word-wrap: break-word;">(无)</td>
<td>String</td>
<td>
- 可以设置和传递任意 Kafka 的配置项。后缀名必须匹配在 <a
href="https://kafka.apache.org/documentation/#configuration">Kafka 配置文档</a>
中定义的配置键。Flink 将移除 "properties." 配置键前缀并将变换后的配置键和值传入底层的 Kafka 客户端。例如,你可以通过
<code>'properties.allow.auto.create.topics' = 'false'</code> 来禁用 topic
的自动创建。但是某些配置项不支持进行配置,因为 Flink 会覆盖这些配置,例如 <code>'key.deserializer'</code> 和
<code>'value.deserializer'</code>。
+ 可以设置和传递任意 Kafka 的配置项。后缀名必须匹配在 <a
href="https://kafka.apache.org/documentation/#configuration">Kafka 配置文档</a>
中定义的配置键。Flink 将移除 "properties." 配置键前缀并将变换后的配置键和值传入底层的 Kafka 客户端。例如,你可以通过
<code>'properties.allow.auto.create.topics' = 'false'</code> 来禁用 topic
的自动创建。<code>'auto.offset.reset'</code> 属性用于配置 source 如何处理 Kafka 中不存在的初始化起始
offset。它独立于
<code>'scan.startup.mode'</code>。由于这两个选项控制不同阶段,因此可以有意地将它们配置为不同的值。但是某些配置项不支持进行配置,因为
Flink 会覆盖这些配置,例如 <code>'key.deserializer'</code> 和 <code>'va [...]
</td>
</tr>
<tr>
diff --git a/docs/content.zh/docs/connectors/table/upsert-kafka.md
b/docs/content.zh/docs/connectors/table/upsert-kafka.md
index bacaae52..251b746c 100644
--- a/docs/content.zh/docs/connectors/table/upsert-kafka.md
+++ b/docs/content.zh/docs/connectors/table/upsert-kafka.md
@@ -136,7 +136,7 @@ of all available metadata fields.
<td>
该选项可以传递任意的 Kafka 参数。选项的后缀名必须匹配定义在 <a
href="https://kafka.apache.org/documentation/#configuration">Kafka
参数文档</a>中的参数名。
Flink 会自动移除 选项名中的 "properties." 前缀,并将转换后的键名以及值传入 KafkaClient。
例如,你可以通过 <code>'properties.allow.auto.create.topics' = 'false'</code>
- 来禁止自动创建 topic。 但是,某些选项,例如<code>'auto.offset.reset'</code>
是不允许通过该方式传递参数,因为 Flink 会重写这些参数的值。
+ 来禁止自动创建 topic。<code>'auto.offset.reset'</code> 属性用于配置 source 如何处理
Kafka 中不存在的初始化起始 offset。它独立于
<code>'scan.startup.mode'</code>。由于这两个选项控制不同阶段,因此可以有意地将它们配置为不同的值。某些其他配置项可能不受支持,因为
Flink 会覆盖它们。
</td>
</tr>
<tr>
diff --git a/docs/content/docs/connectors/datastream/kafka.md
b/docs/content/docs/connectors/datastream/kafka.md
index 60db801a..20290297 100644
--- a/docs/content/docs/connectors/datastream/kafka.md
+++ b/docs/content/docs/connectors/datastream/kafka.md
@@ -236,10 +236,11 @@ For configurations of KafkaConsumer, you can refer to
<a href="http://kafka.apache.org/documentation/#consumerconfigs">Apache Kafka
documentation</a>
for more details.
-Please note that the following keys will be overridden by the builder even if
-it is configured:
-- ```auto.offset.reset.strategy``` is overridden by
```OffsetsInitializer#getAutoOffsetResetStrategy()```
- for the starting offsets
+Please note that the following keys will be set by the builder:
+- ```auto.offset.reset``` is set from
```OffsetsInitializer#getAutoOffsetResetStrategy()```
+ for the starting offsets unless it is explicitly configured. The initializer
selects the initial
+ position, while an explicitly configured ```auto.offset.reset``` controls
what happens if an
+ initialized position later becomes unavailable.
- ```partition.discovery.interval.ms``` is overridden to -1 when
```setBounded(OffsetsInitializer)``` has been invoked
diff --git a/docs/content/docs/connectors/table/kafka.md
b/docs/content/docs/connectors/table/kafka.md
index d73235f4..51e61a6a 100644
--- a/docs/content/docs/connectors/table/kafka.md
+++ b/docs/content/docs/connectors/table/kafka.md
@@ -239,7 +239,7 @@ Connector Options
<td style="word-wrap: break-word;">(none)</td>
<td>String</td>
<td>
- This can set and pass arbitrary Kafka configurations. Suffix names
must match the configuration key defined in <a
href="https://kafka.apache.org/documentation/#configuration">Kafka
Configuration documentation</a>. Flink will remove the "properties." key prefix
and pass the transformed key and values to the underlying KafkaClient. For
example, you can disable automatic topic creation via
<code>'properties.allow.auto.create.topics' = 'false'</code>. But there are
some configuratio [...]
+ This can set and pass arbitrary Kafka configurations. Suffix names
must match the configuration key defined in <a
href="https://kafka.apache.org/documentation/#configuration">Kafka
Configuration documentation</a>. Flink will remove the "properties." key prefix
and pass the transformed key and values to the underlying KafkaClient. For
example, you can disable automatic topic creation via
<code>'properties.allow.auto.create.topics' = 'false'</code>. The
<code>'auto.offset.reset'</ [...]
</td>
</tr>
<tr>
diff --git a/docs/content/docs/connectors/table/upsert-kafka.md
b/docs/content/docs/connectors/table/upsert-kafka.md
index db75309a..974b242d 100644
--- a/docs/content/docs/connectors/table/upsert-kafka.md
+++ b/docs/content/docs/connectors/table/upsert-kafka.md
@@ -144,7 +144,7 @@ Connector Options
<td style="word-wrap: break-word;">(none)</td>
<td>String</td>
<td>
- This can set and pass arbitrary Kafka configurations. Suffix names
must match the configuration key defined in <a
href="https://kafka.apache.org/documentation/#configuration">Kafka
Configuration documentation</a>. Flink will remove the "properties." key prefix
and pass the transformed key and values to the underlying KafkaClient. For
example, you can disable automatic topic creation via
<code>'properties.allow.auto.create.topics' = 'false'</code>. But there are
some configuratio [...]
+ This can set and pass arbitrary Kafka configurations. Suffix names
must match the configuration key defined in <a
href="https://kafka.apache.org/documentation/#configuration">Kafka
Configuration documentation</a>. Flink will remove the "properties." key prefix
and pass the transformed key and values to the underlying KafkaClient. For
example, you can disable automatic topic creation via
<code>'properties.allow.auto.create.topics' = 'false'</code>. The
<code>'auto.offset.reset'</ [...]
</td>
</tr>
<tr>
diff --git
a/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
b/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
index 3341f09d..478ac02a 100644
---
a/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
+++
b/flink-connector-kafka/archunit-violations/c0d94764-76a0-4c50-b617-70b1754c4612
@@ -2,7 +2,7 @@ Class <org.apache.flink.connector.kafka.sink.KafkaSink>
implements interface <or
Class
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator$PartitionChange>
is annotated with <org.apache.flink.annotation.VisibleForTesting> in
(KafkaSourceEnumerator.java:0)
Class
<org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumerator$PartitionOffsetsRetrieverImpl>
is annotated with <org.apache.flink.annotation.VisibleForTesting> in
(KafkaSourceEnumerator.java:0)
Constructor
<org.apache.flink.connector.kafka.dynamic.source.enumerator.DynamicKafkaSourceEnumerator.<init>(org.apache.flink.connector.kafka.dynamic.source.enumerator.subscriber.KafkaStreamSubscriber,
org.apache.flink.connector.kafka.dynamic.metadata.KafkaMetadataService,
org.apache.flink.api.connector.source.SplitEnumeratorContext,
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer,
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInit
[...]
-Constructor
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.<init>(org.apache.flink.api.connector.source.SourceReaderContext,
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema,
java.util.Properties)> calls constructor
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.<init>(int)>
in (DynamicKafkaSourceReader.java:114)
+Constructor
<org.apache.flink.connector.kafka.dynamic.source.reader.DynamicKafkaSourceReader.<init>(org.apache.flink.api.connector.source.SourceReaderContext,
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema,
java.util.Properties,
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer)>
calls constructor
<org.apache.flink.streaming.runtime.io.MultipleFuturesAvailabilityHelper.<init>(int)>
in (DynamicKafkaSourceReader. [...]
Constructor
<org.apache.flink.connector.kafka.sink.KafkaWriterState.<init>(java.lang.String,
int, int, org.apache.flink.connector.kafka.sink.internal.TransactionOwnership,
java.util.Collection)> is annotated with
<org.apache.flink.annotation.VisibleForTesting> in (KafkaWriterState.java:0)
Constructor
<org.apache.flink.streaming.connectors.kafka.table.DynamicKafkaRecordSerializationSchema.<init>(java.util.List,
java.util.regex.Pattern,
org.apache.flink.connector.kafka.sink.KafkaPartitioner,
org.apache.flink.api.common.serialization.SerializationSchema,
org.apache.flink.api.common.serialization.SerializationSchema,
[Lorg.apache.flink.table.data.RowData$FieldGetter;,
[Lorg.apache.flink.table.data.RowData$FieldGetter;, boolean, [I, boolean)> has
parameter of type <[Lorg.apach [...]
Field
<org.apache.flink.connector.kafka.dynamic.source.metrics.KafkaClusterMetricGroupManager.metricGroups>
has generic type <java.util.Map<java.lang.String,
org.apache.flink.runtime.metrics.groups.AbstractMetricGroup>> with type
argument depending on
<org.apache.flink.runtime.metrics.groups.AbstractMetricGroup> in
(KafkaClusterMetricGroupManager.java:0)
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSource.java
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSource.java
index be24f686..0b3609c1 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSource.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSource.java
@@ -126,6 +126,10 @@ public class DynamicKafkaSource<T>
return boundedness;
}
+ Properties getProperties() {
+ return properties;
+ }
+
/**
* Create the {@link DynamicKafkaSourceReader}.
*
@@ -136,7 +140,8 @@ public class DynamicKafkaSource<T>
@Override
public SourceReader<T, DynamicKafkaSourceSplit> createReader(
SourceReaderContext readerContext) {
- return new DynamicKafkaSourceReader<>(readerContext,
deserializationSchema, properties);
+ return new DynamicKafkaSourceReader<>(
+ readerContext, deserializationSchema, properties,
startingOffsetsInitializer);
}
/**
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceBuilder.java
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceBuilder.java
index 9b7c31ba..2bea803f 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceBuilder.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceBuilder.java
@@ -24,6 +24,7 @@ import
org.apache.flink.connector.kafka.dynamic.metadata.KafkaMetadataService;
import
org.apache.flink.connector.kafka.dynamic.source.enumerator.subscriber.KafkaStreamSetSubscriber;
import
org.apache.flink.connector.kafka.dynamic.source.enumerator.subscriber.KafkaStreamSubscriber;
import
org.apache.flink.connector.kafka.dynamic.source.enumerator.subscriber.StreamPatternSubscriber;
+import org.apache.flink.connector.kafka.source.KafkaPropertiesUtil;
import org.apache.flink.connector.kafka.source.KafkaSourceOptions;
import
org.apache.flink.connector.kafka.source.enumerator.initializer.NoStoppingOffsetsInitializer;
import
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
@@ -33,6 +34,7 @@ import org.apache.flink.util.Preconditions;
import org.apache.commons.lang3.RandomStringUtils;
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.common.serialization.ByteArrayDeserializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -254,10 +256,6 @@ public class DynamicKafkaSourceBuilder<T> {
maybeOverride(KafkaSourceOptions.COMMIT_OFFSETS_ON_CHECKPOINT.key(), "false",
false);
}
maybeOverride(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false",
false);
- maybeOverride(
- ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
-
startingOffsetsInitializer.getAutoOffsetResetStrategy().name().toLowerCase(),
- true);
// If the source is bounded, do not run periodic partition discovery.
maybeOverride(
@@ -317,6 +315,30 @@ public class DynamicKafkaSourceBuilder<T> {
String.format(
"Property %s is required when offset commit is
enabled",
ConsumerConfig.GROUP_ID_CONFIG));
+
+ warnIfOffsetResetStrategyOpposesStartingOffsetsInitializer();
+ }
+
+ private void warnIfOffsetResetStrategyOpposesStartingOffsetsInitializer() {
+ String configuredOffsetReset =
props.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG);
+ if (configuredOffsetReset == null) {
+ return;
+ }
+
+ OffsetResetStrategy configuredOffsetResetStrategy =
+ KafkaPropertiesUtil.getResetStrategy(configuredOffsetReset);
+ if (KafkaPropertiesUtil.hasOpposingOffsetResetStrategies(
+ configuredOffsetResetStrategy, startingOffsetsInitializer)) {
+ logger.warn(
+ "Configured {}={} differs from the {} strategy derived
from the starting "
+ + "offsets initializer. The source will use the
initializer for "
+ + "startup, but Kafka may reset to {} if an
initialized offset "
+ + "becomes unavailable.",
+ ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ configuredOffsetReset,
+
startingOffsetsInitializer.getAutoOffsetResetStrategy().name().toLowerCase(),
+ configuredOffsetReset);
+ }
}
private boolean offsetCommitEnabledManually() {
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumerator.java
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumerator.java
index d2cd5eea..32bd3b51 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumerator.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/enumerator/DynamicKafkaSourceEnumerator.java
@@ -45,6 +45,7 @@ import org.apache.flink.util.Preconditions;
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.common.KafkaException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -615,12 +616,12 @@ public class DynamicKafkaSourceEnumerator
KafkaPropertiesUtil.copyProperties(properties, consumerProps);
DynamicKafkaSourceOptions.removeRemovedClusterRetentionOption(consumerProps);
KafkaPropertiesUtil.setClientIdPrefix(consumerProps, kafkaClusterId);
+ OffsetResetStrategy effectiveOffsetResetStrategy =
+ KafkaPropertiesUtil.resolveAutoOffsetResetStrategy(
+ properties, fetchedProperties,
effectiveStartingOffsetsInitializer);
consumerProps.setProperty(
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
- effectiveStartingOffsetsInitializer
- .getAutoOffsetResetStrategy()
- .name()
- .toLowerCase());
+ effectiveOffsetResetStrategy.name().toLowerCase());
KafkaSourceEnumerator enumerator =
new KafkaSourceEnumerator(
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReader.java
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReader.java
index 3f4d0469..2a6bd266 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReader.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReader.java
@@ -88,6 +88,7 @@ public class DynamicKafkaSourceReader<T> implements
SourceReader<T, DynamicKafka
private final KafkaRecordDeserializationSchema<T> deserializationSchema;
private final Properties properties;
+ private final OffsetsInitializer startingOffsetsInitializer;
private final MetricGroup dynamicKafkaSourceMetricGroup;
private final Gauge<Integer> kafkaClusterCount;
private final AtomicInteger activeSplitCount;
@@ -114,10 +115,19 @@ public class DynamicKafkaSourceReader<T> implements
SourceReader<T, DynamicKafka
SourceReaderContext readerContext,
KafkaRecordDeserializationSchema<T> deserializationSchema,
Properties properties) {
+ this(readerContext, deserializationSchema, properties,
OffsetsInitializer.earliest());
+ }
+
+ public DynamicKafkaSourceReader(
+ SourceReaderContext readerContext,
+ KafkaRecordDeserializationSchema<T> deserializationSchema,
+ Properties properties,
+ OffsetsInitializer startingOffsetsInitializer) {
this.readerContext = readerContext;
this.clusterReaderMap = new TreeMap<>();
this.deserializationSchema = deserializationSchema;
this.properties = properties;
+ this.startingOffsetsInitializer = startingOffsetsInitializer;
this.kafkaClusterCount = clusterReaderMap::size;
this.activeSplitCount = new AtomicInteger();
this.dynamicKafkaSourceMetricGroup =
@@ -284,16 +294,20 @@ public class DynamicKafkaSourceReader<T> implements
SourceReader<T, DynamicKafka
Properties clusterProperties = new Properties();
KafkaPropertiesUtil.copyProperties(
clusterMetadataMapEntry.getValue().getProperties(),
clusterProperties);
- OffsetsInitializer startingOffsetsInitializer =
+ OffsetsInitializer clusterStartingOffsetsInitializer =
clusterMetadataMapEntry.getValue().getStartingOffsetsInitializer();
- if (startingOffsetsInitializer != null) {
- clusterProperties.setProperty(
- ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
- startingOffsetsInitializer
- .getAutoOffsetResetStrategy()
- .name()
- .toLowerCase());
- }
+ OffsetsInitializer effectiveStartingOffsetsInitializer =
+ clusterStartingOffsetsInitializer != null
+ ? clusterStartingOffsetsInitializer
+ : startingOffsetsInitializer;
+ clusterProperties.setProperty(
+ ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ KafkaPropertiesUtil.resolveAutoOffsetResetStrategy(
+ properties,
+ clusterProperties,
+ effectiveStartingOffsetsInitializer)
+ .name()
+ .toLowerCase());
newClustersProperties.put(clusterMetadataMapEntry.getKey(),
clusterProperties);
}
}
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaPropertiesUtil.java
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaPropertiesUtil.java
index 0e29576c..9a13df06 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaPropertiesUtil.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaPropertiesUtil.java
@@ -19,10 +19,17 @@
package org.apache.flink.connector.kafka.source;
import org.apache.flink.annotation.Internal;
+import
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
+
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import javax.annotation.Nonnull;
+import java.util.Arrays;
+import java.util.Locale;
import java.util.Properties;
+import java.util.stream.Collectors;
/** Utility class for modify Kafka properties. */
@Internal
@@ -36,6 +43,57 @@ public class KafkaPropertiesUtil {
}
}
+ /** Resolves an explicit cluster or global reset strategy before the
initializer default. */
+ public static OffsetResetStrategy resolveAutoOffsetResetStrategy(
+ @Nonnull Properties globalProperties,
+ @Nonnull Properties clusterProperties,
+ @Nonnull OffsetsInitializer startingOffsetsInitializer) {
+ return getResetStrategy(
+ clusterProperties.getProperty(
+ ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ globalProperties.getProperty(
+ ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+
startingOffsetsInitializer.getAutoOffsetResetStrategy().name())));
+ }
+
+ /** Parses the configured auto offset reset strategy. */
+ public static OffsetResetStrategy getResetStrategy(@Nonnull String
offsetResetConfig) {
+ return Arrays.stream(OffsetResetStrategy.values())
+ .filter(
+ offsetResetStrategy ->
+ offsetResetStrategy
+ .name()
+
.equals(offsetResetConfig.toUpperCase(Locale.ROOT)))
+ .findAny()
+ .orElseThrow(
+ () ->
+ new IllegalArgumentException(
+ String.format(
+ "%s can not be set to %s.
Valid values: [%s]",
+
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ offsetResetConfig,
+
Arrays.stream(OffsetResetStrategy.values())
+ .map(Enum::name)
+
.map(String::toLowerCase)
+
.collect(Collectors.joining(",")))));
+ }
+
+ /** Returns whether the configured strategy opposes a positional
initializer strategy. */
+ public static boolean hasOpposingOffsetResetStrategies(
+ @Nonnull OffsetResetStrategy configuredResetStrategy,
+ @Nonnull OffsetsInitializer startingOffsetsInitializer) {
+ OffsetResetStrategy initializerResetStrategy =
+ startingOffsetsInitializer.getAutoOffsetResetStrategy();
+ return isPositionalResetStrategy(configuredResetStrategy)
+ && isPositionalResetStrategy(initializerResetStrategy)
+ && configuredResetStrategy != initializerResetStrategy;
+ }
+
+ private static boolean isPositionalResetStrategy(OffsetResetStrategy
resetStrategy) {
+ return resetStrategy == OffsetResetStrategy.EARLIEST
+ || resetStrategy == OffsetResetStrategy.LATEST;
+ }
+
/**
* client.id is used for Kafka server side logging, see
*
https://docs.confluent.io/platform/current/installation/configuration/consumer-configs.html#consumerconfigs_client.id
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilder.java
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilder.java
index 0709afe0..4167f385 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilder.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilder.java
@@ -29,6 +29,7 @@ import
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDe
import org.apache.flink.util.function.SerializableSupplier;
import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.ByteArrayDeserializer;
import org.apache.kafka.common.serialization.Deserializer;
@@ -382,9 +383,9 @@ public class KafkaSourceBuilder<OUT> {
* created.
*
* <ul>
- * <li><code>auto.offset.reset.strategy</code> is overridden by {@link
- * OffsetsInitializer#getAutoOffsetResetStrategy()} for the starting
offsets, which is by
- * default {@link OffsetsInitializer#earliest()}.
+ * <li><code>auto.offset.reset</code> is set from {@link
+ * OffsetsInitializer#getAutoOffsetResetStrategy()} for the starting
offsets unless
+ * explicitly configured by the user.
* <li><code>partition.discovery.interval.ms</code> is overridden to -1
when {@link
* #setBounded(OffsetsInitializer)} has been invoked.
* </ul>
@@ -406,9 +407,9 @@ public class KafkaSourceBuilder<OUT> {
* created.
*
* <ul>
- * <li><code>auto.offset.reset.strategy</code> is overridden by {@link
- * OffsetsInitializer#getAutoOffsetResetStrategy()} for the starting
offsets, which is by
- * default {@link OffsetsInitializer#earliest()}.
+ * <li><code>auto.offset.reset</code> is set from {@link
+ * OffsetsInitializer#getAutoOffsetResetStrategy()} for the starting
offsets unless
+ * explicitly configured by the user.
* <li><code>partition.discovery.interval.ms</code> is overridden to -1
when {@link
* #setBounded(OffsetsInitializer)} has been invoked.
* <li><code>client.id</code> is overridden to the
"client.id.prefix-RANDOM_LONG", or
@@ -468,10 +469,32 @@ public class KafkaSourceBuilder<OUT> {
maybeOverride(KafkaSourceOptions.COMMIT_OFFSETS_ON_CHECKPOINT.key(), "false",
false);
}
maybeOverride(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false",
false);
+ String configuredOffsetReset =
props.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG);
+ if (configuredOffsetReset != null) {
+ OffsetResetStrategy configuredOffsetResetStrategy =
+
KafkaPropertiesUtil.getResetStrategy(configuredOffsetReset);
+ String normalizedOffsetReset =
configuredOffsetResetStrategy.name().toLowerCase();
+ props.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
normalizedOffsetReset);
+ if (KafkaPropertiesUtil.hasOpposingOffsetResetStrategies(
+ configuredOffsetResetStrategy,
startingOffsetsInitializer)) {
+ LOG.warn(
+ "Configured {}={} differs from the {} strategy derived
from the starting "
+ + "offsets initializer. The source will use
the initializer for "
+ + "startup, but Kafka may reset to {} if an
initialized offset "
+ + "becomes unavailable.",
+ ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ normalizedOffsetReset,
+ startingOffsetsInitializer
+ .getAutoOffsetResetStrategy()
+ .name()
+ .toLowerCase(),
+ normalizedOffsetReset);
+ }
+ }
maybeOverride(
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
startingOffsetsInitializer.getAutoOffsetResetStrategy().name().toLowerCase(),
- true);
+ false);
// If the source is bounded, do not run periodic partition discovery.
if (boundedness == Boundedness.BOUNDED) {
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaTableSource.java
b/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaTableSource.java
index da303b7a..5e007b3e 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaTableSource.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaTableSource.java
@@ -26,6 +26,7 @@ import org.apache.flink.api.connector.source.Boundedness;
import org.apache.flink.connector.kafka.dynamic.metadata.KafkaMetadataService;
import org.apache.flink.connector.kafka.dynamic.source.DynamicKafkaSource;
import
org.apache.flink.connector.kafka.dynamic.source.DynamicKafkaSourceBuilder;
+import org.apache.flink.connector.kafka.source.KafkaPropertiesUtil;
import
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
import org.apache.flink.streaming.api.datastream.DataStream;
@@ -66,7 +67,6 @@ import java.util.HashMap;
import java.util.HashSet;
import java.util.LinkedHashMap;
import java.util.List;
-import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
@@ -443,6 +443,7 @@ public class DynamicKafkaTableSource
DeserializationSchema<RowData> keyDeserialization,
DeserializationSchema<RowData> valueDeserialization,
TypeInformation<RowData> producedTypeInfo) {
+
KafkaConnectorOptionsUtil.validateAndNormalizeAutoOffsetResetStrategy(properties);
final KafkaRecordDeserializationSchema<RowData> kafkaDeserializer =
createKafkaDeserializationSchema(
@@ -462,31 +463,7 @@ public class DynamicKafkaTableSource
.setDeserializer(kafkaDeserializer)
.setProperties(properties);
- switch (startupMode) {
- case EARLIEST:
-
dynamicKafkaSourceBuilder.setStartingOffsets(OffsetsInitializer.earliest());
- break;
- case LATEST:
-
dynamicKafkaSourceBuilder.setStartingOffsets(OffsetsInitializer.latest());
- break;
- case GROUP_OFFSETS:
- String offsetResetConfig =
- properties.getProperty(
- ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
- OffsetResetStrategy.NONE.name());
- OffsetResetStrategy offsetResetStrategy =
getResetStrategy(offsetResetConfig);
- dynamicKafkaSourceBuilder.setStartingOffsets(
-
OffsetsInitializer.committedOffsets(offsetResetStrategy));
- break;
- case SPECIFIC_OFFSETS:
- dynamicKafkaSourceBuilder.setStartingOffsets(
- OffsetsInitializer.offsets(specificStartupOffsets));
- break;
- case TIMESTAMP:
- dynamicKafkaSourceBuilder.setStartingOffsets(
- OffsetsInitializer.timestamp(startupTimestampMillis));
- break;
- }
+
dynamicKafkaSourceBuilder.setStartingOffsets(getStartingOffsetsInitializer());
switch (boundedMode) {
case UNBOUNDED:
@@ -510,21 +487,36 @@ public class DynamicKafkaTableSource
return dynamicKafkaSourceBuilder.build();
}
- private OffsetResetStrategy getResetStrategy(String offsetResetConfig) {
- return Arrays.stream(OffsetResetStrategy.values())
- .filter(ors ->
ors.name().equals(offsetResetConfig.toUpperCase(Locale.ROOT)))
- .findAny()
- .orElseThrow(
- () ->
- new IllegalArgumentException(
- String.format(
- "%s can not be set to %s.
Valid values: [%s]",
-
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
- offsetResetConfig,
-
Arrays.stream(OffsetResetStrategy.values())
- .map(Enum::name)
-
.map(String::toLowerCase)
-
.collect(Collectors.joining(",")))));
+ private OffsetsInitializer getStartingOffsetsInitializer() {
+ final OffsetsInitializer startingOffsetsInitializer;
+ switch (startupMode) {
+ case EARLIEST:
+ startingOffsetsInitializer = OffsetsInitializer.earliest();
+ break;
+ case LATEST:
+ startingOffsetsInitializer = OffsetsInitializer.latest();
+ break;
+ case GROUP_OFFSETS:
+ String offsetResetConfig =
+ properties.getProperty(
+ ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ OffsetResetStrategy.NONE.name());
+ OffsetResetStrategy offsetResetStrategy =
+
KafkaPropertiesUtil.getResetStrategy(offsetResetConfig);
+ startingOffsetsInitializer =
+
OffsetsInitializer.committedOffsets(offsetResetStrategy);
+ break;
+ case SPECIFIC_OFFSETS:
+ startingOffsetsInitializer =
OffsetsInitializer.offsets(specificStartupOffsets);
+ break;
+ case TIMESTAMP:
+ startingOffsetsInitializer =
OffsetsInitializer.timestamp(startupTimestampMillis);
+ break;
+ default:
+ throw new IllegalStateException("Unsupported startup mode: " +
startupMode);
+ }
+
+ return startingOffsetsInitializer;
}
private KafkaRecordDeserializationSchema<RowData>
createKafkaDeserializationSchema(
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaConnectorOptionsUtil.java
b/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaConnectorOptionsUtil.java
index ecaa3091..fff07609 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaConnectorOptionsUtil.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaConnectorOptionsUtil.java
@@ -44,15 +44,19 @@ import org.apache.flink.util.FlinkException;
import org.apache.flink.util.InstantiationUtil;
import org.apache.flink.util.Preconditions;
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.common.TopicPartition;
import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.Optional;
import java.util.Properties;
import java.util.regex.Pattern;
+import java.util.stream.Collectors;
import java.util.stream.IntStream;
import static
org.apache.flink.streaming.connectors.kafka.table.KafkaConnectorOptions.DELIVERY_GUARANTEE;
@@ -115,6 +119,34 @@ class KafkaConnectorOptionsUtil {
validateSinkPartitioner(tableOptions);
}
+ static void validateAndNormalizeAutoOffsetResetStrategy(Properties
properties) {
+ String resetStrategy =
properties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG);
+ if (resetStrategy == null) {
+ return;
+ }
+ String normalizedResetStrategy =
resetStrategy.toLowerCase(Locale.ROOT);
+
+ boolean valid =
+ Arrays.stream(OffsetResetStrategy.values())
+ .anyMatch(
+ strategy ->
+ strategy.name()
+ .toLowerCase(Locale.ROOT)
+
.equals(normalizedResetStrategy));
+ if (!valid) {
+ throw new IllegalArgumentException(
+ String.format(
+ "%s can not be set to %s. Valid values: [%s]",
+ ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ resetStrategy,
+ Arrays.stream(OffsetResetStrategy.values())
+ .map(Enum::name)
+ .map(value ->
value.toLowerCase(Locale.ROOT))
+ .collect(Collectors.joining(","))));
+ }
+ properties.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
normalizedResetStrategy);
+ }
+
public static void validateTopic(ReadableConfig tableOptions) {
Optional<List<String>> topic = tableOptions.getOptional(TOPIC);
Optional<String> pattern = tableOptions.getOptional(TOPIC_PATTERN);
diff --git
a/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicSource.java
b/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicSource.java
index 39f7b014..c0f7ba2c 100644
---
a/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicSource.java
+++
b/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicSource.java
@@ -23,6 +23,7 @@ import
org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.serialization.DeserializationSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.connector.source.Boundedness;
+import org.apache.flink.connector.kafka.source.KafkaPropertiesUtil;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.KafkaSourceBuilder;
import
org.apache.flink.connector.kafka.source.enumerator.initializer.NoStoppingOffsetsInitializer;
@@ -67,7 +68,6 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.LinkedHashMap;
import java.util.List;
-import java.util.Locale;
import java.util.Map;
import java.util.Objects;
import java.util.Optional;
@@ -431,6 +431,7 @@ public class KafkaDynamicSource
DeserializationSchema<RowData> keyDeserialization,
DeserializationSchema<RowData> valueDeserialization,
TypeInformation<RowData> producedTypeInfo) {
+
KafkaConnectorOptionsUtil.validateAndNormalizeAutoOffsetResetStrategy(properties);
final KafkaRecordDeserializationSchema<RowData> kafkaDeserializer =
createKafkaDeserializationSchema(
@@ -444,79 +445,71 @@ public class KafkaDynamicSource
kafkaSourceBuilder.setTopicPattern(topicPattern);
}
- switch (startupMode) {
- case EARLIEST:
-
kafkaSourceBuilder.setStartingOffsets(OffsetsInitializer.earliest());
+ kafkaSourceBuilder.setStartingOffsets(getStartingOffsetsInitializer());
+
+ switch (boundedMode) {
+ case UNBOUNDED:
+ kafkaSourceBuilder.setUnbounded(new
NoStoppingOffsetsInitializer());
break;
case LATEST:
-
kafkaSourceBuilder.setStartingOffsets(OffsetsInitializer.latest());
+ kafkaSourceBuilder.setBounded(OffsetsInitializer.latest());
break;
case GROUP_OFFSETS:
- String offsetResetConfig =
- properties.getProperty(
- ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
- OffsetResetStrategy.NONE.name());
- OffsetResetStrategy offsetResetStrategy =
getResetStrategy(offsetResetConfig);
- kafkaSourceBuilder.setStartingOffsets(
-
OffsetsInitializer.committedOffsets(offsetResetStrategy));
+
kafkaSourceBuilder.setBounded(OffsetsInitializer.committedOffsets());
break;
case SPECIFIC_OFFSETS:
Map<TopicPartition, Long> offsets = new HashMap<>();
- specificStartupOffsets.forEach(
+ specificBoundedOffsets.forEach(
(tp, offset) ->
offsets.put(
new TopicPartition(tp.topic(),
tp.partition()), offset));
-
kafkaSourceBuilder.setStartingOffsets(OffsetsInitializer.offsets(offsets));
+
kafkaSourceBuilder.setBounded(OffsetsInitializer.offsets(offsets));
break;
case TIMESTAMP:
- kafkaSourceBuilder.setStartingOffsets(
- OffsetsInitializer.timestamp(startupTimestampMillis));
+
kafkaSourceBuilder.setBounded(OffsetsInitializer.timestamp(boundedTimestampMillis));
break;
}
- switch (boundedMode) {
- case UNBOUNDED:
- kafkaSourceBuilder.setUnbounded(new
NoStoppingOffsetsInitializer());
+
kafkaSourceBuilder.setProperties(properties).setDeserializer(kafkaDeserializer);
+
+ return kafkaSourceBuilder.build();
+ }
+
+ private OffsetsInitializer getStartingOffsetsInitializer() {
+ final OffsetsInitializer startingOffsetsInitializer;
+ switch (startupMode) {
+ case EARLIEST:
+ startingOffsetsInitializer = OffsetsInitializer.earliest();
break;
case LATEST:
- kafkaSourceBuilder.setBounded(OffsetsInitializer.latest());
+ startingOffsetsInitializer = OffsetsInitializer.latest();
break;
case GROUP_OFFSETS:
-
kafkaSourceBuilder.setBounded(OffsetsInitializer.committedOffsets());
+ String offsetResetConfig =
+ properties.getProperty(
+ ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ OffsetResetStrategy.NONE.name());
+ OffsetResetStrategy offsetResetStrategy =
+
KafkaPropertiesUtil.getResetStrategy(offsetResetConfig);
+ startingOffsetsInitializer =
+
OffsetsInitializer.committedOffsets(offsetResetStrategy);
break;
case SPECIFIC_OFFSETS:
Map<TopicPartition, Long> offsets = new HashMap<>();
- specificBoundedOffsets.forEach(
+ specificStartupOffsets.forEach(
(tp, offset) ->
offsets.put(
new TopicPartition(tp.topic(),
tp.partition()), offset));
-
kafkaSourceBuilder.setBounded(OffsetsInitializer.offsets(offsets));
+ startingOffsetsInitializer =
OffsetsInitializer.offsets(offsets);
break;
case TIMESTAMP:
-
kafkaSourceBuilder.setBounded(OffsetsInitializer.timestamp(boundedTimestampMillis));
+ startingOffsetsInitializer =
OffsetsInitializer.timestamp(startupTimestampMillis);
break;
+ default:
+ throw new IllegalStateException("Unsupported startup mode: " +
startupMode);
}
-
kafkaSourceBuilder.setProperties(properties).setDeserializer(kafkaDeserializer);
-
- return kafkaSourceBuilder.build();
- }
-
- private OffsetResetStrategy getResetStrategy(String offsetResetConfig) {
- return Arrays.stream(OffsetResetStrategy.values())
- .filter(ors ->
ors.name().equals(offsetResetConfig.toUpperCase(Locale.ROOT)))
- .findAny()
- .orElseThrow(
- () ->
- new IllegalArgumentException(
- String.format(
- "%s can not be set to %s.
Valid values: [%s]",
-
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
- offsetResetConfig,
-
Arrays.stream(OffsetResetStrategy.values())
- .map(Enum::name)
-
.map(String::toLowerCase)
-
.collect(Collectors.joining(",")))));
+ return startingOffsetsInitializer;
}
private KafkaRecordDeserializationSchema<RowData>
createKafkaDeserializationSchema(
diff --git
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceBuilderTest.java
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceBuilderTest.java
new file mode 100644
index 00000000..3b71a01f
--- /dev/null
+++
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/DynamicKafkaSourceBuilderTest.java
@@ -0,0 +1,102 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.connector.kafka.dynamic.source;
+
+import org.apache.flink.connector.kafka.dynamic.metadata.KafkaMetadataService;
+import org.apache.flink.connector.kafka.dynamic.metadata.KafkaStream;
+import
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
+import
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
+
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.common.serialization.IntegerDeserializer;
+import org.junit.jupiter.api.Test;
+
+import java.util.Collection;
+import java.util.Collections;
+import java.util.Map;
+import java.util.Set;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests for {@link DynamicKafkaSourceBuilder}. */
+class DynamicKafkaSourceBuilderTest {
+
+ @Test
+ void testAutoOffsetResetIsNotMaterializedWhenAbsent() {
+ assertThat(
+ baseBuilder()
+ .build()
+ .getProperties()
+
.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
+ .isNull();
+ }
+
+ @Test
+ void testAutoOffsetResetUsesExplicitProperty() {
+ assertThat(
+ baseBuilder()
+
.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "none")
+ .build()
+ .getProperties()
+
.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
+ .isEqualTo("none");
+ }
+
+ @Test
+ void testAutoOffsetResetExplicitPropertyOverridesInitializerStrategy() {
+ assertThat(
+ baseBuilder()
+
.setStartingOffsets(OffsetsInitializer.latest())
+
.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "none")
+ .build()
+ .getProperties()
+
.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
+ .isEqualTo("none");
+ }
+
+ private DynamicKafkaSourceBuilder<Integer> baseBuilder() {
+ return DynamicKafkaSource.<Integer>builder()
+ .setStreamIds(Collections.singleton("stream-1"))
+ .setKafkaMetadataService(NoOpKafkaMetadataService.INSTANCE)
+ .setDeserializer(
+
KafkaRecordDeserializationSchema.valueOnly(IntegerDeserializer.class));
+ }
+
+ private enum NoOpKafkaMetadataService implements KafkaMetadataService {
+ INSTANCE;
+
+ @Override
+ public Set<KafkaStream> getAllStreams() {
+ return Collections.emptySet();
+ }
+
+ @Override
+ public Map<String, KafkaStream> describeStreams(Collection<String>
streamIds) {
+ return Collections.emptyMap();
+ }
+
+ @Override
+ public boolean isClusterActive(String kafkaClusterId) {
+ return true;
+ }
+
+ @Override
+ public void close() {}
+ }
+}
diff --git
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReaderTest.java
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReaderTest.java
index c7a4ec2a..ebffa3e9 100644
---
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReaderTest.java
+++
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/dynamic/source/reader/DynamicKafkaSourceReaderTest.java
@@ -32,6 +32,7 @@ import
org.apache.flink.connector.kafka.dynamic.source.DynamicKafkaSourceOptions
import org.apache.flink.connector.kafka.dynamic.source.MetadataUpdateEvent;
import
org.apache.flink.connector.kafka.dynamic.source.metrics.KafkaClusterMetricGroup;
import
org.apache.flink.connector.kafka.dynamic.source.split.DynamicKafkaSourceSplit;
+import
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import
org.apache.flink.connector.kafka.source.metrics.KafkaSourceReaderMetrics;
import org.apache.flink.connector.kafka.source.reader.KafkaSourceReader;
import
org.apache.flink.connector.kafka.source.reader.deserializer.KafkaRecordDeserializationSchema;
@@ -471,7 +472,8 @@ public class DynamicKafkaSourceReaderTest extends
SourceReaderTestBase<DynamicKa
return new DynamicKafkaSourceReader<>(
context,
KafkaRecordDeserializationSchema.valueOnly(IntegerDeserializer.class),
- properties);
+ properties,
+ OffsetsInitializer.earliest());
}
private DynamicKafkaSourceReader<Integer>
createReaderWithoutStartWithRemovedClusterRetention(
@@ -483,7 +485,8 @@ public class DynamicKafkaSourceReaderTest extends
SourceReaderTestBase<DynamicKa
return new DynamicKafkaSourceReader<>(
context,
KafkaRecordDeserializationSchema.valueOnly(IntegerDeserializer.class),
- properties);
+ properties,
+ OffsetsInitializer.earliest());
}
private SourceReader<Integer, DynamicKafkaSourceSplit> startReader(
@@ -627,7 +630,8 @@ class DynamicKafkaSourceReaderPauseResumeTest {
return new DynamicKafkaSourceReader<>(
new TestingReaderContext(),
KafkaRecordDeserializationSchema.valueOnly(IntegerDeserializer.class),
- properties);
+ properties,
+ OffsetsInitializer.earliest());
}
private static DynamicKafkaSourceSplit createSplit(
diff --git
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaPropertiesUtilTest.java
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaPropertiesUtilTest.java
new file mode 100644
index 00000000..7363f2c0
--- /dev/null
+++
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaPropertiesUtilTest.java
@@ -0,0 +1,107 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.connector.kafka.source;
+
+import
org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
+
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.OffsetResetStrategy;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.Arguments;
+import org.junit.jupiter.params.provider.MethodSource;
+
+import java.util.Arrays;
+import java.util.Properties;
+import java.util.stream.Stream;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Tests for {@link KafkaPropertiesUtil}. */
+class KafkaPropertiesUtilTest {
+
+ @Test
+ void testUsesInitializerStrategyWhenResetPropertiesAreAbsent() {
+ assertThat(
+ KafkaPropertiesUtil.resolveAutoOffsetResetStrategy(
+ new Properties(), new Properties(),
OffsetsInitializer.earliest()))
+ .isEqualTo(OffsetResetStrategy.EARLIEST);
+
+ assertThat(
+ KafkaPropertiesUtil.resolveAutoOffsetResetStrategy(
+ new Properties(), new Properties(),
OffsetsInitializer.latest()))
+ .isEqualTo(OffsetResetStrategy.LATEST);
+ }
+
+ @Test
+ void testClusterResetPropertyOverridesInitializerStrategy() {
+ assertThat(
+ KafkaPropertiesUtil.resolveAutoOffsetResetStrategy(
+ new Properties(),
+ resetProperties("none"),
+ OffsetsInitializer.earliest()))
+ .isEqualTo(OffsetResetStrategy.NONE);
+ }
+
+ @Test
+ void testClusterResetPropertyOverridesGlobalAndInitializerStrategies() {
+ assertThat(
+ KafkaPropertiesUtil.resolveAutoOffsetResetStrategy(
+ resetProperties("none"),
+ resetProperties("earliest"),
+ OffsetsInitializer.latest()))
+ .isEqualTo(OffsetResetStrategy.EARLIEST);
+ }
+
+ @ParameterizedTest
+ @MethodSource("allOffsetResetStrategyPairs")
+ void testDetectsOnlyOpposingPositionalResetStrategies(
+ OffsetResetStrategy configuredResetStrategy,
+ OffsetResetStrategy initializerResetStrategy) {
+ boolean expected =
+ (configuredResetStrategy == OffsetResetStrategy.EARLIEST
+ && initializerResetStrategy ==
OffsetResetStrategy.LATEST)
+ || (configuredResetStrategy ==
OffsetResetStrategy.LATEST
+ && initializerResetStrategy ==
OffsetResetStrategy.EARLIEST);
+
+ assertThat(
+ KafkaPropertiesUtil.hasOpposingOffsetResetStrategies(
+ configuredResetStrategy,
+
OffsetsInitializer.committedOffsets(initializerResetStrategy)))
+ .isEqualTo(expected);
+ }
+
+ private static Stream<Arguments> allOffsetResetStrategyPairs() {
+ return Arrays.stream(OffsetResetStrategy.values())
+ .flatMap(
+ configuredResetStrategy ->
+ Arrays.stream(OffsetResetStrategy.values())
+ .map(
+ initializerResetStrategy ->
+ Arguments.of(
+
configuredResetStrategy,
+
initializerResetStrategy)));
+ }
+
+ private static Properties resetProperties(String resetStrategy) {
+ Properties properties = new Properties();
+ properties.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
resetStrategy);
+ return properties;
+ }
+}
diff --git
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilderTest.java
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilderTest.java
index 921c2b56..bb7d71c0 100644
---
a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilderTest.java
+++
b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilderTest.java
@@ -85,6 +85,42 @@ public class KafkaSourceBuilderTest {
.isFalse();
}
+ @Test
+ public void testAutoOffsetResetDefaultsToInitializerStrategy() {
+
assertThat(getAutoOffsetResetStrategy(getBasicBuilder().build())).isEqualTo("earliest");
+ }
+
+ @Test
+ public void testAutoOffsetResetUsesExplicitProperty() {
+ KafkaSource<String> kafkaSource =
+ getBasicBuilder()
+ .setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
"none")
+ .build();
+
+ assertThat(getAutoOffsetResetStrategy(kafkaSource)).isEqualTo("none");
+ }
+
+ @Test
+ public void testAutoOffsetResetNormalizesExplicitProperty() {
+ KafkaSource<String> kafkaSource =
+ getBasicBuilder()
+ .setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
"EARLIEST")
+ .build();
+
+
assertThat(getAutoOffsetResetStrategy(kafkaSource)).isEqualTo("earliest");
+ }
+
+ @Test
+ public void
testAutoOffsetResetExplicitPropertyOverridesInitializerStrategy() {
+ KafkaSource<String> kafkaSource =
+ getBasicBuilder()
+ .setStartingOffsets(OffsetsInitializer.latest())
+ .setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
"none")
+ .build();
+
+ assertThat(getAutoOffsetResetStrategy(kafkaSource)).isEqualTo("none");
+ }
+
@Test
public void testEnableCommitOnCheckpointWithoutGroupId() {
assertThatThrownBy(
@@ -245,6 +281,15 @@ public class KafkaSourceBuilderTest {
KafkaRecordDeserializationSchema.valueOnly(StringDeserializer.class));
}
+ private String getAutoOffsetResetStrategy(KafkaSource<?> kafkaSource) {
+ return kafkaSource
+ .getConfiguration()
+ .get(
+
ConfigOptions.key(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG)
+ .stringType()
+ .noDefaultValue());
+ }
+
private static class ExampleCustomSubscriber implements KafkaSubscriber {
@Override
diff --git
a/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaTableFactoryTest.java
b/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaTableFactoryTest.java
index 3bb5deae..dda08c07 100644
---
a/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaTableFactoryTest.java
+++
b/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaTableFactoryTest.java
@@ -28,6 +28,7 @@ import org.apache.flink.table.catalog.ResolvedSchema;
import org.apache.flink.table.connector.source.DynamicTableSource;
import org.apache.flink.table.factories.TestFormatFactory;
+import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.junit.jupiter.api.Test;
import java.util.Arrays;
@@ -97,6 +98,18 @@ class DynamicKafkaTableFactoryTest {
.isEqualTo("60000");
}
+ @Test
+ void testTableSourcePreservesConfiguredOffsetResetStrategy() {
+ final Map<String, String> options = getSingleClusterSourceOptions();
+ options.put("properties." + ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
"none");
+
+ final DynamicKafkaTableSource tableSource =
+ (DynamicKafkaTableSource) createTableSource(SCHEMA, options);
+
+
assertThat(tableSource.properties.getProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
+ .isEqualTo("none");
+ }
+
private static Map<String, String> getSingleClusterSourceOptions() {
Map<String, String> tableOptions = new HashMap<>();
// Dynamic Kafka specific options.
diff --git
a/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicTableFactoryTest.java
b/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicTableFactoryTest.java
index 5f309344..63cc1b24 100644
---
a/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicTableFactoryTest.java
+++
b/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicTableFactoryTest.java
@@ -88,6 +88,7 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
+import java.util.Locale;
import java.util.Map;
import java.util.Optional;
import java.util.Properties;
@@ -468,21 +469,72 @@ class KafkaDynamicTableFactoryTest {
}
@ParameterizedTest
- @ValueSource(strings = {"none", "earliest", "latest"})
+ @ValueSource(strings = {"none", "earliest", "latest", "EARLIEST"})
@NullSource
public void testTableSourceSetOffsetReset(final String strategyName) {
testSetOffsetResetForStartFromGroupOffsets(strategyName);
}
- @Test
- void testTableSourceSetOffsetResetWithException() {
- String errorStrategy = "errorStrategy";
- assertThatThrownBy(() -> testTableSourceSetOffsetReset(errorStrategy))
+ @ParameterizedTest
+ @ValueSource(
+ strings = {
+ "earliest-offset",
+ "latest-offset",
+ "specific-offsets",
+ "timestamp",
+ "group-offsets"
+ })
+ void testTableSourceSetOffsetResetForEveryStartupMode(String startupMode) {
+ final Map<String, String> modifiedOptions =
+ getModifiedOptions(
+ getBasicSourceOptions(),
+ options -> {
+ options.put("scan.startup.mode", startupMode);
+ options.put(
+ PROPERTIES_PREFIX +
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ "none");
+ if (!"specific-offsets".equals(startupMode)) {
+
options.remove("scan.startup.specific-offsets");
+ }
+ if ("timestamp".equals(startupMode)) {
+ options.put("scan.startup.timestamp-millis",
"1000");
+ }
+ });
+
+
assertThat(getTableSourceAutoOffsetReset(modifiedOptions)).isEqualTo("none");
+ }
+
+ @ParameterizedTest
+ @ValueSource(
+ strings = {
+ "earliest-offset",
+ "latest-offset",
+ "specific-offsets",
+ "timestamp",
+ "group-offsets"
+ })
+ void testTableSourceRejectsInvalidOffsetResetForEveryStartupMode(String
startupMode) {
+ final Map<String, String> modifiedOptions =
+ getModifiedOptions(
+ getBasicSourceOptions(),
+ options -> {
+ options.put("scan.startup.mode", startupMode);
+ options.put(
+ PROPERTIES_PREFIX +
ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
+ "errorStrategy");
+ if (!"specific-offsets".equals(startupMode)) {
+
options.remove("scan.startup.specific-offsets");
+ }
+ if ("timestamp".equals(startupMode)) {
+ options.put("scan.startup.timestamp-millis",
"1000");
+ }
+ });
+
+ assertThatThrownBy(() ->
getTableSourceAutoOffsetReset(modifiedOptions))
.isInstanceOf(IllegalArgumentException.class)
.hasMessage(
- String.format(
- "%s can not be set to %s. Valid values:
[latest,earliest,none]",
- ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,
errorStrategy));
+ "auto.offset.reset can not be set to errorStrategy. "
+ + "Valid values: [latest,earliest,none]");
}
private void testSetOffsetResetForStartFromGroupOffsets(String value) {
@@ -499,6 +551,15 @@ class KafkaDynamicTableFactoryTest {
value);
});
final DynamicTableSource tableSource = createTableSource(SCHEMA,
modifiedOptions);
+ assertThat(getTableSourceAutoOffsetReset(tableSource))
+ .isEqualTo(value == null ? "none" :
value.toLowerCase(Locale.ROOT));
+ }
+
+ private String getTableSourceAutoOffsetReset(Map<String, String> options) {
+ return getTableSourceAutoOffsetReset(createTableSource(SCHEMA,
options));
+ }
+
+ private String getTableSourceAutoOffsetReset(DynamicTableSource
tableSource) {
assertThat(tableSource).isInstanceOf(KafkaDynamicSource.class);
ScanTableSource.ScanRuntimeProvider provider =
((KafkaDynamicSource) tableSource)
@@ -508,13 +569,7 @@ class KafkaDynamicTableFactoryTest {
final Configuration configuration =
KafkaSourceTestUtils.getKafkaSourceConfiguration(kafkaSource);
- if (value == null) {
-
assertThat(configuration.toMap().get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
- .isEqualTo("none");
- } else {
-
assertThat(configuration.toMap().get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG))
- .isEqualTo(value);
- }
+ return
configuration.toMap().get(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG);
}
@Test