You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Copy file name to clipboardExpand all lines: flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/KafkaSourceEnumerator.java
+1-1
Original file line number
Diff line number
Diff line change
@@ -230,7 +230,7 @@ public void close() {
230
230
* @return Set of subscribed {@link TopicPartition}s
Copy file name to clipboardExpand all lines: flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/KafkaSubscriber.java
+4-1
Original file line number
Diff line number
Diff line change
@@ -25,6 +25,7 @@
25
25
26
26
importjava.io.Serializable;
27
27
importjava.util.List;
28
+
importjava.util.Properties;
28
29
importjava.util.Set;
29
30
importjava.util.regex.Pattern;
30
31
@@ -51,9 +52,11 @@ public interface KafkaSubscriber extends Serializable {
51
52
* Get a set of subscribed {@link TopicPartition}s.
52
53
*
53
54
* @param adminClient The admin client used to retrieve subscribed topic partitions.
55
+
* @param properties The properties for the configuration.
54
56
* @return A set of subscribed {@link TopicPartition}s
Copy file name to clipboardExpand all lines: flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/KafkaSubscriberUtils.java
Copy file name to clipboardExpand all lines: flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/PartitionSetSubscriber.java
+4-2
Original file line number
Diff line number
Diff line change
@@ -30,6 +30,7 @@
30
30
importjava.util.HashSet;
31
31
importjava.util.Map;
32
32
importjava.util.Optional;
33
+
importjava.util.Properties;
33
34
importjava.util.Set;
34
35
importjava.util.stream.Collectors;
35
36
@@ -46,15 +47,16 @@ class PartitionSetSubscriber implements KafkaSubscriber, KafkaDatasetIdentifierP
Copy file name to clipboardExpand all lines: flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicListSubscriber.java
Copy file name to clipboardExpand all lines: flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/TopicPatternSubscriber.java
+4-2
Original file line number
Diff line number
Diff line change
@@ -31,6 +31,7 @@
31
31
importjava.util.HashSet;
32
32
importjava.util.Map;
33
33
importjava.util.Optional;
34
+
importjava.util.Properties;
34
35
importjava.util.Set;
35
36
importjava.util.regex.Pattern;
36
37
@@ -47,10 +48,11 @@ class TopicPatternSubscriber implements KafkaSubscriber, KafkaDatasetIdentifierP
Copy file name to clipboardExpand all lines: flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/enumerator/subscriber/KafkaSubscriberTest.java
+8-5
Original file line number
Diff line number
Diff line change
@@ -33,6 +33,7 @@
33
33
importjava.util.Collections;
34
34
importjava.util.HashSet;
35
35
importjava.util.List;
36
+
importjava.util.Properties;
36
37
importjava.util.Set;
37
38
importjava.util.regex.Pattern;
38
39
@@ -46,13 +47,15 @@ public class KafkaSubscriberTest {
0 commit comments