kafka FileStreamSourceConnector write an avro file to topic with key field

Viewed 425

I want to use kafka FileStreamSourceConnector to write a local avro file into a topic.

My connector config looks like this:

curl -i -X PUT -H  "Content-Type:application/json" http://localhost:8083/connectors/file_source_connector/config \
            -d '{
            "connector.class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
            "value.converter.schema.registry.url": "http://schema-registry:8081",
            "topic": "my_topic",
            "file": "/data/log.avsc",
            "format.include.keys": "true",
            "source.auto.offset.reset": "earliest",
            "tasks.max": "1",
            "value.converter.schemas.enable": "true",
            "value.converter": "io.confluent.connect.avro.AvroConverter",
            "key.converter": "org.apache.kafka.connect.storage.StringConverter"
          }'

Then when I print out the topic, the key fields are null.

Updated on 2021-03-29:

After watching this video Twelve Days of SMT - Day 2: ValueToKey and ExtractField from Robin, I applied SMT to my connector config:

curl -i -X PUT -H  "Content-Type:application/json" http://localhost:8083/connectors/file_source_connector_02/config \
            -d '{
            "connector.class": "org.apache.kafka.connect.file.FileStreamSourceConnector",
            "value.converter.schema.registry.url": "http://schema-registry:8081",
            "topic": "my_topic",
            "file": "/data/log.avsc",
            "tasks.max": "1",
            "value.converter": "io.confluent.connect.avro.AvroConverter",
            "key.converter": "org.apache.kafka.connect.storage.StringConverter",
            "transforms": "ValueToKey, ExtractField",
            "transforms.ValueToKey.type":"org.apache.kafka.connect.transforms.ValueToKey",
            "transforms.ValueToKey.fields":"id",
            "transforms.ExtractField.type":"org.apache.kafka.connect.transforms.ExtractField$Key",
            "transforms.ExtractField.field":"id"
          }'

However, the connector is failed:

Caused by: org.apache.kafka.connect.errors.DataException: Only Struct objects supported for [copying fields from value to key], found: java.lang.String
2 Answers

I would use ValueToKey transformer. In bad case ignorig values and setting random key.

For details look at:ValueToKey

FileStreamSource assumes UTF8 encoded, line delimited files are your input, not binary files such as Avro. Last I checked, format.include.keys is not a valid config for the connector either.

Therefore each consumed event will be a string, and subsequently, transforms that require Structs with field names will not work

You can use the Hoist transform to create a Struct from each "line", but this still will not parse your data to make the ID field accessible to move to the key.

Also, your file is AVSC, which is JSON formatted, not Avro, so I'm not sure what the goal is by using the AvroConverter, or having "schemas.enable": "true". Still, the lines read by the connector are not parsed by converters such that fields are accessible, only serialized when sent to Kafka


My suggestion would be to write some other CLI script using plain producer libraries to parse the file, extract the schema, register that with Schema Registry, build a producer record for each entity in the file, and send them

Related