Github user fhueske commented on a diff in the pull request:
https://github.com/apache/flink/pull/5056#discussion_r152760280
--- Diff:
flink-connectors/flink-connector-kafka-base/src/main/java/org/apache/flink/streaming/connectors/kafka/KafkaTableSource.java
---
@@ -121,6 +130,37 @@ public String getProctimeAttribute() {
return rowtimeAttributeDescriptors;
}
+ /**
+ * Returns a version-specific Kafka consumer with the start position
configured.
+ *
+ * @param topic Kafka topic to consume.
+ * @param properties Properties for the Kafka consumer.
+ * @param deserializationSchema Deserialization schema to use for Kafka
records.
+ * @return The version-specific Kafka consumer
+ */
+ public FlinkKafkaConsumerBase<Row> getKafkaConsumer(
--- End diff --
Should be `protected`
---