Kafka Connect 写入Cassandra时主键uid映射为空的问题求助
解决Kafka Connect Cassandra Sink的主键uid映射为空问题
看起来你遇到的核心问题是Zeek日志的嵌套结构和Kafka Connect配置的语法/路径错误导致Cassandra无法正确提取uid字段。我们一步步来排查修复:
1. 先修复转换器配置的语法错误
你的cassandra-sink.properties里的transforms配置误用了JSON格式的引号和逗号,这不符合Kafka Connect properties文件的键值对语法,必须改成标准格式:
# 转换Zeek字段名以避免错误 transforms=RenameField transforms.RenameField.type=org.apache.kafka.connect.transforms.ReplaceField$Value transforms.RenameField.renames=id.orig_h:id_orig_h,id.orig_p:id_orig_p,id.resp_h:id_resp_h,id.resp_p:id_resp_p,AA:aa,Z:z,RA:ra,TC:tc,TS:ts,RD:rd
注意这里加了$Value后缀,指定转换器作用于消息的value部分(也就是嵌套的dns/http对象内部的字段),否则转换器会尝试修改整个顶层消息结构,达不到预期效果。
2. 修正字段映射路径,指向嵌套的Zeek字段
从你的Kafka日志示例可以看到,所有Zeek字段都被包裹在dns(或http)对象里,比如{"dns": {"ts":1644391084.805351,"uid":"CUXnim2Q50AeUZoZXc",...}}。而你原来的映射规则uid=key是错误的(因为日志的key是null),uid=value会把整个顶层对象作为uid的值,自然无法匹配Cassandra的text类型主键。
你需要明确指定嵌套路径来提取字段,修改后的映射规则如下:
DNS主题映射
topic.dns.zeek.dns.mapping= ts=dns.value.ts,uid=dns.value.uid,id_orig_h=dns.value.id.orig_h,id_orig_p=dns.value.id.orig_p,id_resp_h=dns.value.id.resp_h,id_resp_p=dns.value.id.resp_p,proto=dns.value.proto,trans_id=dns.value.trans_id,rtt=dns.value.rtt,query=dns.value.query,qclass=dns.value.qclass,qclass_name=dns.value.qclass_name,qtype=dns.value.qtype,qtype_name=dns.value.qtype_name,rcode=dns.value.rcode,rcode_name=dns.value.rcode_name,aa=dns.value.AA,tc=dns.value.TC,rd=dns.value.RD,ra=dns.value.RA,z=dns.value.Z,answers=dns.value.answers,rejected=dns.value.rejected
HTTP主题映射
topic.http.zeek.http.mapping= ts=http.value.ts,uid=http.value.uid,id_orig_h=http.value.id.orig_h,id_orig_p=http.value.id.orig_p,id_resp_h=http.value.id.resp_h,id_resp_p=http.value.id.resp_p,trans_depth=http.value.trans_depth,method=http.value.method,host=http.value.host,uri=http.value.uri,version=http.value.version,user_agent=http.value.user_agent,request_body_len=http.value.request_body_len,response_body_len=http.value.response_body_len,status_code=http.value.status_code,status_msg=http.value.status_msg,tags=http.value.tags,resp_fuids=http.value.resp_fuids
3. 验证配置和数据匹配
- 确认Cassandra表的字段类型和Zeek日志字段类型匹配:比如你的
dns表中id_orig_p是double,Zeek里是整数,存入double没问题;uid是text,Zeek的uid是字符串,完全匹配。 - 启动Connect后,可以通过Kafka Connect的REST API(默认
http://localhost:8083/connectors/cassandra-sink/config)查看生效的配置,确保没有语法错误。 - 如果还是有问题,可以临时去掉转换器,先让映射直接用原始字段名(比如
id.orig_h),验证数据能正常写入后再加上字段重命名。
为什么之前的配置不生效?
- 转换器的语法错误导致字段重命名没有生效,Cassandra找不到对应的字段。
- 映射路径没有指向嵌套的
dns/http对象内部,导致Connect尝试从顶层消息体提取uid,结果拿到的是整个对象而非字符串,最终解析为null。
内容的提问来源于stack exchange,提问作者kafkaconnecth8
相关产品推荐
相关产品推荐

