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 <[email protected]>
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);
}