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

Reply via email to