You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Kafka Connect反序列化Divolte的Kafka Avro记录并转JSON至Kinesis

如何用Kafka Connect消费并反序列化Divolte Collector写入Kafka的Avro记录?

我需要把Divolte写入Kafka的Avro序列化记录转换成JSON事件,再导入Kinesis数据流。目前使用AWS Labs Kinesis Streams Sink插件,遇到了两个问题:

  • 最初将Divolte设为裸模式,使用无注册表Avro转换器(Worker配置中已注释),报错not an Avro file,无法正常工作。
  • 切换Divolte为Confluent模式,使用io.confluent.connect.avro.AvroConverter转换器,出现以下错误:

kafka-connect| Caused by: org.apache.kafka.common.errors.SerializationException: Unknown magic byte!


Worker配置

bootstrap.servers=broker:29092

key.converter=io.confluent.connect.avro.AvroConverter
value.converter=io.confluent.connect.avro.AvroConverter
key.converter.schema.registry.url=http://localhost:8081/subjects/Kafka-key/versions/1
value.converter.schema.registry.url=http://localhost:8081/subjects/Kafka-key/versions/1

#key.converter=me.frmr.kafka.connect.RegistrylessAvroConverter
#value.converter=me.frmr.kafka.connect.RegistrylessAvroConverter
#key.converter.schema.path=/opt/kafka/config/DefaultEventRecord.avsc
#value.converter.schema.path=/opt/kafka/config/DefaultEventRecord.avsc

key.converter.schemas.enable=false
value.converter.schemas.enable=false

#internal.value.converter=org.apache.kafka.connect.storage.StringConverter
#internal.key.converter=org.apache.kafka.connect.storage.StringConverter
internal.value.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter=org.apache.kafka.connect.json.JsonConverter
internal.key.converter.schemas.enable=true
internal.value.converter.schemas.enable=true

offset.storage.file.filename=offset.log
schemas.enable=false
plugin.path=/opt/kafka/plugins/

Divolte Collector配置

已将sink设为Confluent模式,并指定confluent_id=1(对应Schema Registry中提交的第一个Schema ID):

divolte {
  global {
    server {
      host = 0.0.0.0
      host = ${?DIVOLTE_HOST}
      port = 8290
      port = ${?DIVOLTE_PORT}
      use_x_forwarded_for = false
      use_x_forwarded_for = ${?DIVOLTE_USE_XFORWARDED_FOR}
      serve_static_resources = true
      serve_static_resources = ${?DIVOLTE_SERVICE_STATIC_RESOURCES}
      debug_requests = false
    }

    mapper {
      buffer_size = 1048576
      threads = 1
      duplicate_memory_size = 1000000
      user_agent_parser {
        type = non_updating
        cache_size = 1000
      }
    }

    kafka {
      enabled = false
      enabled = ${?DIVOLTE_KAFKA_ENABLED}
      threads = 2
      buffer_size = 1048576
      producer = {
        bootstrap.servers = ["localhost:9092"]
        bootstrap.servers = ${?DIVOLTE_KAFKA_BROKER_LIST}
        client.id = divolte.collector
        client.id = ${?DIVOLTE_KAFKA_CLIENT_ID}
        acks = 1
        retries = 0
        compression.type = lz4
        max.in.flight.requests.per.connection = 1

        sasl.jaas.config = ""
        sasl.jaas.config = ${?KAFKA_SASL_JAAS_CONFIG}

        security.protocol = PLAINTEXT
        security.protocol = ${?KAFKA_SECURITY_PROTOCOL}
        sasl.mechanism = GSSAPI
        sasl.kerberos.service.name = kafka
      }
    }
  }

  sources {
    browser1 = {
      type = browser
      event_suffix = event
      party_cookie = _dvp
      party_cookie = ${?DIVOLTE_PARTY_COOKIE}
      party_timeout = 730 days
      party_timeout = ${?DIVOLTE_PARTY_TIMEOUT}
      session_cookie = _dvs
      session_cookie = ${?DIVOLTE_SESSION_COOKIE}
      session_timeout = 30 minutes
      session_timeout = ${?DIVOLTE_SESSION_TIMEOUT}
      cookie_domain = ''
      cookie_domain = ${?DIVOLTE_COOKIE_DOMAIN}

      javascript {
        name = divolte.js
        name = ${?DIVOLTE_JAVASCRIPT_NAME}
        logging = false
        logging = ${?DIVOLTE_JAVASCRIPT_LOGGING}
        debug = false
        debug = ${?DIVOLTE_JAVASCRIPT_DEBUG}
        auto_page_view_event = true
        auto_page_view_event = ${?DIVOLTE_JAVASCRIPT_AUTO_PAGE_VIEW_EVENT}
      }
    }
  }

  sinks {
    kafka1 = {
      type = kafka
      mode = confluent
      topic = clickstream
      topic = ${?DIVOLTE_KAFKA_TOPIC}
    }
  }

  mappings {
    a_mapping = {
    //schema_file = /opt/divolte/conf/DefaultEventRecord.avsc
    //mapping_script_file = schema-mapping.groovy
    confluent_id = 1
    sources = [browser1]
    sinks = [kafka1]
    }
  }
}

Schema Registry相关操作

使用Divolte默认的DefaultEventRecord.avsc Schema,内容如下:

{"schema":"{\"namespace\":\"io.divolte.record\",\"type\":\"record\",\"name\":\"DefaultEventRecord\",\"fields\":[{\"name\":\"detectedDuplicate\",\"type\":\"boolean\"},{\"name\":\"detectedCorruption\",\"type\":\"boolean\"},{\"name\":\"firstInSession\",\"type\":\"boolean\"},{\"name\":\"timestamp\",\"type\":\"long\"},{\"name\":\"clientTimestamp\",\"type\":\"long\"},{\"name\":\"remoteHost\",\"type\":\"string\"},{\"name\":\"referer\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"location\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"viewportPixelWidth\",\"type\":[\"null\",\"int\"],\"default\":null},{\"name\":\"viewportPixelHeight\",\"type\":[\"null\",\"int\"],\"default\":null},{\"name\":\"screenPixelWidth\",\"type\":[\"null\",\"int\"],\"default\":null},{\"name\":\"screenPixelHeight\",\"type\":[\"null\",\"int\"],\"default\":null},{\"name\":\"partyId\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"sessionId\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"pageViewId\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"eventType\",\"type\":\"string\",\"default\":\"unknown\"},{\"name\":\"userAgentString\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"userAgentName\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"userAgentFamily\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"userAgentVendor\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"userAgentType\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"userAgentVersion\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"userAgentDeviceCategory\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"userAgentOsFamily\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"userAgentOsVersion\",\"type\":[\"null\",\"string\"],\"default\":null},{\"name\":\"userAgentOsVendor\",\"type\":[\"null\",\"string\"],\"default\":null}]}"} 

通过以下命令提交到Schema Registry:

curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" --data @divolte-schema-v1.avsc http://localhost:8081/subjects/Kafka-key/versions

并将Schema Registry的版本URL配置到了Worker中。


疑问点

  • 我的Schema文件是否正确?
  • 是否使用了适配Divolte Avro序列化的转换器?
  • Divolte的Confluent模式配置是否正确,confluent_id是否有效?(我设置的ID=1是Schema Registry启动后提交第一个Schema返回的ID)

内容的提问来源于stack exchange,提问作者Adam Prior

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.25 12:57:03