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); }