sjvanrossum commented on code in PR #40181:
URL: https://github.com/apache/beam/pull/40181#discussion_r4056470064
##########
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"
Review Comment:
Nit: Google auth library docs typically use the name "Application Default
Credentials" or in rare cases "Google Application Default Credentials" since
the provider can be used for Google APIs in general.
##########
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();
Review Comment:
Nit: This doesn't match the name used to configure the transform.
Application default credentials is annoyingly long to spell so I'm not opposed
to ADC, but I'd consider keeping it the same for consistency and clarity to
users who might be unfamiliar with the acronym.
##########
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:
In most cases the partition set for a topic will match the range of integers
from 0 to N, but users may arbitrarily retire individual partitions.
Examples:
- Create a partition per entity with a non-sequential integer id (e.g., a
partition per player in a game).
- Retire a partition after the last usable offset is used to produce a
record.
Replacing the number of partitions property with a topic partition list
property would spare adding it later with more boilerplate to validate the
different combinations of configuration properties.
--
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]