|
24 | 24 | import java.math.MathContext; |
25 | 25 | import java.time.Duration; |
26 | 26 | import java.util.Collections; |
27 | | -import java.util.HashMap; |
28 | 27 | import java.util.List; |
29 | 28 | import java.util.Map; |
30 | 29 | import java.util.Optional; |
@@ -211,15 +210,6 @@ private ReadFromKafkaDoFn( |
211 | 210 | transform.getAdminFactoryFn(); |
212 | 211 | final SerializableFunction<Map<String, Object>, Consumer<byte[], byte[]>> consumerFactoryFn = |
213 | 212 | transform.getConsumerFactoryFn(); |
214 | | - final @Nullable Map<String, Object> offsetConsumerConfigOverrides = |
215 | | - transform.getOffsetConsumerConfig(); |
216 | | - final Map<String, Object> offsetConsumerConfig; |
217 | | - if (offsetConsumerConfigOverrides == null) { |
218 | | - offsetConsumerConfig = transform.getConsumerConfig(); |
219 | | - } else { |
220 | | - offsetConsumerConfig = new HashMap<>(transform.getConsumerConfig()); |
221 | | - offsetConsumerConfig.putAll(offsetConsumerConfigOverrides); |
222 | | - } |
223 | 213 | this.consumerConfig = transform.getConsumerConfig(); |
224 | 214 | this.keyDeserializerProvider = |
225 | 215 | Preconditions.checkArgumentNotNull(transform.getKeyDeserializerProvider()); |
@@ -270,7 +260,7 @@ public KafkaLatestOffsetEstimator load( |
270 | 260 | sourceDescriptor); |
271 | 261 | final Map<String, Object> config = |
272 | 262 | KafkaIOUtils.overrideBootstrapServersConfig( |
273 | | - offsetConsumerConfig, sourceDescriptor); |
| 263 | + consumerConfig, sourceDescriptor); |
274 | 264 | final Admin admin = adminFactoryFn.apply(config); |
275 | 265 | return new KafkaLatestOffsetEstimator( |
276 | 266 | admin, sourceDescriptor.getTopicPartition()); |
|
0 commit comments