【问题标题】:Kafka connect - > S3 Parquet file byteArreyKafka 连接 -> S3 Parquet 文件 byteArrey
【发布时间】:2022-07-07 23:42:37
【问题描述】:

我正在尝试使用 Kafka-connect 来消耗 Kafka 的消息并将它们写入 s3 parquet 文件。 所以我写了一个简单的生产者,用 byte[]

生成消息
Properties propertiesAWS = new Properties();
    propertiesAWS.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "myKafka:9092");
    propertiesAWS.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, LongSerializer.class.getName());
    propertiesAWS.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName());


    KafkaProducer<Long, byte[]> producer = new KafkaProducer<Long, byte[]>(propertiesAWS);
    Random rng = new Random();

    for (int i = 0; i < 100; i++) {
        try {
            Thread.sleep(1000);
            Headers headers = new RecordHeaders();
            headers.add(new RecordHeader("header1", "header1".getBytes()));
            headers.add(new RecordHeader("header2", "header2".getBytes()));
            ProducerRecord<Long, byte[]> recordOut = new ProducerRecord<Long, byte[]>
                    ("s3.test.topic", 1, rng.nextLong(), new byte[]{1, 2, 3}, headers);
            producer.send(recordOut);
        } catch (Exception e) {
            System.out.println(e);

        }
    }

我的 kafka 连接配置是:

{
"name": "test_2_s3",
"config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "aws.access.key.id": "XXXXXXX",
    "aws.secret.access.key": "XXXXXXXX",
    "s3.region": "eu-central-1",
    "flush.size": "5",
    "rotate.schedule.interval.ms": "10000",
    "timezone": "UTC",
    "tasks.max": "1",
    "topics": "s3.test.topic",
    "parquet.codec": "gzip",
    "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
    "partitioner.class": "io.confluent.connect.storage.partitioner.DefaultPartitioner",
    "storage.class": "io.confluent.connect.s3.storage.S3Storage",
    "s3.bucket.name": "test-phase1",
    "key.converter": "org.apache.kafka.connect.converters.LongConverter",
    "value.converter": "org.apache.kafka.connect.converters.ByteArrayConverter",
    "behavior.on.null.values": "ignore",
    "store.kafka.headers": "true"
}

这是我得到的错误:

原因:java.lang.IllegalArgumentException:Avro 模式必须是记录。 在 org.apache.parquet.avro.AvroSchemaConverter.convert(AvroSchemaConverter.java:124)

我的错误在哪里? 我是否需要使用 Avro,即使我只想写 byteArr + 一些 Kafka 头文件? 如何配置要写入镶木地板的 kafka 标头? 谢谢

【问题讨论】:

    标签: java amazon-s3 apache-kafka parquet apache-kafka-connect


    【解决方案1】:

    ParquetFormat 编写器需要架构为 type: record 的 Avro 数据,是的。

    不过,您可以使用这样的架构。

    {
        "namespace": "com.stackoverflow.example",
        "name": "ByteWrapper",
        "type": "record",
        "fields": [{
            "name": "data",
            "type": "bytes"
        }]
    }
    

    配置要写入 parquet 的 kafka 标头

    没有。 Kafka header 只是字节,没有指定的序列化格式。

    如果您只想将字节写入 S3,请改用 ByteArrayFormat

    【讨论】:

      猜你喜欢
      • 2017-10-08
      • 2021-10-06
      • 1970-01-01
      • 2020-01-18
      • 2020-02-26
      • 2019-10-28
      • 2021-02-25
      • 2022-10-14
      • 2021-10-12
      相关资源
      最近更新 更多