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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 12:37:34