Abacn commented on code in PR #40181:
URL: https://github.com/apache/beam/pull/40181#discussion_r4064428536
##########
sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaReadSchemaTransformConfiguration.java:
##########
@@ -198,6 +198,21 @@ public static Builder builder() {
@Nullable
public abstract Boolean getRedistributeByRecordKey();
+ @SchemaFieldDescription(
+ "Whether to use Google Cloud Platform Application Default Credentials
(ADC) for"
+ + " authenticating with a Google Managed Kafka cluster.")
+ @SchemaFieldNumber("17")
+ @Nullable
+ public abstract Boolean getWithGcpAdc();
+
+ @SchemaFieldDescription(
+ "The number of partitions to read from the Kafka topic. If specified,
partitions"
+ + " 0 to numPartitions-1 will be assigned statically without
querying Kafka at pipeline"
+ + " construction time.")
+ @SchemaFieldNumber("18")
+ @Nullable
+ public abstract Integer getNumPartitions();
Review Comment:
It's true there is a risk number of partitions evolving won't get honored,
but in fact it's the same limitation for a typical setup
- on Dataflow streaming runner, it uses UnboundedSource, which resolves
partitions at once, at initial split that happens at pipeline expansion time
- on Dataflow portable runner, it uses SDF implementation. Initial split
happens at pipeline execution time, but after that the partitions to consume is
still fixed
There is an `.withDynamicRead(...)` option that support dynamic partitions,
led to SDK implementation and requires Dataflow runner v2.
One of the motivation of exposing this option is to be able to use KafkaIO
via SchemaTransform to connect to a GMK on Dataflow Streaming runner, therefore
`.withDynamicRead(...)` isn't available here
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]