如何用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
相关产品推荐
相关产品推荐

