|
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; |
@@ -201,15 +200,6 @@ private ReadFromKafkaDoFn( |
201 | 200 | transform.getAdminFactoryFn(); |
202 | 201 | final SerializableFunction<Map<String, Object>, Consumer<byte[], byte[]>> consumerFactoryFn = |
203 | 202 | transform.getConsumerFactoryFn(); |
204 | | - final @Nullable Map<String, Object> offsetConsumerConfigOverrides = |
205 | | - transform.getOffsetConsumerConfig(); |
206 | | - final Map<String, Object> offsetConsumerConfig; |
207 | | - if (offsetConsumerConfigOverrides == null) { |
208 | | - offsetConsumerConfig = transform.getConsumerConfig(); |
209 | | - } else { |
210 | | - offsetConsumerConfig = new HashMap<>(transform.getConsumerConfig()); |
211 | | - offsetConsumerConfig.putAll(offsetConsumerConfigOverrides); |
212 | | - } |
213 | 203 | this.consumerConfig = transform.getConsumerConfig(); |
214 | 204 | this.keyDeserializerProvider = |
215 | 205 | Preconditions.checkArgumentNotNull(transform.getKeyDeserializerProvider()); |
@@ -260,7 +250,7 @@ public KafkaLatestOffsetEstimator load( |
260 | 250 | sourceDescriptor); |
261 | 251 | final Map<String, Object> config = |
262 | 252 | KafkaIOUtils.overrideBootstrapServersConfig( |
263 | | - offsetConsumerConfig, sourceDescriptor); |
| 253 | + consumerConfig, sourceDescriptor); |
264 | 254 | final Admin admin = adminFactoryFn.apply(config); |
265 | 255 | return new KafkaLatestOffsetEstimator( |
266 | 256 | admin, sourceDescriptor.getTopicPartition()); |
|
0 commit comments