@AutoService(value=org.apache.beam.sdk.schemas.transforms.SchemaTransformProvider.class) public class KafkaReadSchemaTransformProvider extends org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider<KafkaReadSchemaTransformConfiguration>
| Modifier and Type | Class and Description |
|---|---|
static class |
KafkaReadSchemaTransformProvider.ErrorFn |
| Modifier and Type | Field and Description |
|---|---|
static org.apache.beam.sdk.values.TupleTag<org.apache.beam.sdk.values.Row> |
ERROR_TAG |
static org.apache.beam.sdk.values.TupleTag<org.apache.beam.sdk.values.Row> |
OUTPUT_TAG |
| Constructor and Description |
|---|
KafkaReadSchemaTransformProvider() |
| Modifier and Type | Method and Description |
|---|---|
protected java.lang.Class<KafkaReadSchemaTransformConfiguration> |
configurationClass() |
protected org.apache.beam.sdk.schemas.transforms.SchemaTransform |
from(KafkaReadSchemaTransformConfiguration configuration) |
static org.apache.beam.sdk.transforms.SerializableFunction<byte[],org.apache.beam.sdk.values.Row> |
getRawBytesToRowFunction(org.apache.beam.sdk.schemas.Schema rawSchema) |
java.lang.String |
identifier() |
java.util.List<java.lang.String> |
inputCollectionNames() |
java.util.List<java.lang.String> |
outputCollectionNames() |
public static final org.apache.beam.sdk.values.TupleTag<org.apache.beam.sdk.values.Row> OUTPUT_TAG
public static final org.apache.beam.sdk.values.TupleTag<org.apache.beam.sdk.values.Row> ERROR_TAG
protected java.lang.Class<KafkaReadSchemaTransformConfiguration> configurationClass()
configurationClass in class org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider<KafkaReadSchemaTransformConfiguration>protected org.apache.beam.sdk.schemas.transforms.SchemaTransform from(KafkaReadSchemaTransformConfiguration configuration)
from in class org.apache.beam.sdk.schemas.transforms.TypedSchemaTransformProvider<KafkaReadSchemaTransformConfiguration>public static org.apache.beam.sdk.transforms.SerializableFunction<byte[],org.apache.beam.sdk.values.Row> getRawBytesToRowFunction(org.apache.beam.sdk.schemas.Schema rawSchema)
public java.lang.String identifier()
public java.util.List<java.lang.String> inputCollectionNames()
public java.util.List<java.lang.String> outputCollectionNames()