Divolte与Kafka Docker环境下Avro Schema缺失及解析错误求助
解决Divolte Kafka消息转Avro的Schema问题
我之前也碰到过一模一样的问题!Divolte的Kafka消息并不是纯Avro格式——它前面加了自定义的前缀(就是你错误日志里的divolte::3::0),而且很多人容易忽略它内置的Schema获取途径,下面一步步帮你解决:
1. 获取正确的Avro Schema
Divolte自带了两种获取Schema的方式,都能拿到你需要的标准Schema:
- 通过内置API查询:
Divolte默认在8290端口提供Schema注册接口,如果你已经把容器的8290端口映射到宿主机,直接在宿主机执行:
接口会返回所有Divolte使用的Schema,其中curl http://localhost:8290/api/v1/schemasdivolte-collection对应的就是默认收集事件的Avro Schema,复制出来即可使用。 - 从容器内直接读取Schema文件:
进入Divolte容器,默认的Schema文件存放在/opt/divolte/divolte-collector-*/conf/schemas/目录下,执行以下命令查看:
这个docker exec -it <你的divolte容器名称> bash cat /opt/divolte/divolte-collector-*/conf/schemas/collection.avsccollection.avsc就是你需要的标准Avro Schema文件内容。
2. 处理Divolte的消息前缀
拿到Schema直接解析还是会报错,因为Kafka消息开头的divolte::x::x前缀不属于Avro数据部分,必须先剥离:
在StreamSets的流水线里,Kafka源阶段之后添加一个Script Evaluator组件,用Groovy脚本处理消息:
// 获取原始消息字符串 def rawMsg = record.value.valueAsString // 找到换行符位置(Divolte前缀与Avro数据用换行分隔) def payloadStart = rawMsg.indexOf('\n') + 1 // 截取纯Avro二进制内容并替换原消息 record.value = rawMsg.substring(payloadStart).getBytes()
处理完成后再用Avro解析器处理消息,就能正常解析了。
3. 自定义事件的特殊情况
如果你的Divolte配置了自定义收集事件,记得去conf/schemas目录下找对应的自定义Schema文件,或者通过API接口里的其他Schema条目获取,操作步骤和上面一致。
内容的提问来源于stack exchange,提问作者Vadim
相关产品推荐
相关产品推荐

