This is an automated email from the ASF dual-hosted git repository.

yhu pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 67a39a50df3 make BeamKafkaTable.createKafkaRead to be protected, and 
we can override it appending options to KafkaIO.Read (#29051)
67a39a50df3 is described below

commit 67a39a50df3f677d9d1e8c8519f80b222c1101fb
Author: gabry.wu <gabr...@apache.org>
AuthorDate: Fri Oct 20 05:56:41 2023 +0800

    make BeamKafkaTable.createKafkaRead to be protected, and we can override it 
appending options to KafkaIO.Read (#29051)
---
 .../beam/sdk/extensions/sql/meta/provider/kafka/BeamKafkaTable.java     | 2 +-
 .../beam/sdk/extensions/sql/meta/provider/kafka/KafkaTestTable.java     | 2 +-
 2 files changed, 2 insertions(+), 2 deletions(-)

diff --git 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/BeamKafkaTable.java
 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/BeamKafkaTable.java
index f1ec20831a4..ab1817f6d75 100644
--- 
a/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/BeamKafkaTable.java
+++ 
b/sdks/java/extensions/sql/src/main/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/BeamKafkaTable.java
@@ -110,7 +110,7 @@ public abstract class BeamKafkaTable extends 
SchemaBaseBeamTable {
         .setRowSchema(getSchema());
   }
 
-  KafkaIO.Read<byte[], byte[]> createKafkaRead() {
+  protected KafkaIO.Read<byte[], byte[]> createKafkaRead() {
     KafkaIO.Read<byte[], byte[]> kafkaRead;
     if (topics != null) {
       kafkaRead =
diff --git 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/KafkaTestTable.java
 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/KafkaTestTable.java
index 44b4dbe21ac..158b0345bd8 100644
--- 
a/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/KafkaTestTable.java
+++ 
b/sdks/java/extensions/sql/src/test/java/org/apache/beam/sdk/extensions/sql/meta/provider/kafka/KafkaTestTable.java
@@ -61,7 +61,7 @@ public class KafkaTestTable extends BeamKafkaTable {
   }
 
   @Override
-  KafkaIO.Read<byte[], byte[]> createKafkaRead() {
+  protected KafkaIO.Read<byte[], byte[]> createKafkaRead() {
     return super.createKafkaRead().withConsumerFactoryFn(this::mkMockConsumer);
   }
 

Reply via email to