☰
Arkime Kafka 写入插件(Kafka Write Plugin)完全指南:配置、构建与原理
2026/9/28 2:59:25 网站建设 项目流程
  • 网络安全
  • 网络
  • 后端
  • 数据可视化

【免费下载链接】arkime

Arkime is an open source, large scale, full packet capturing, indexing, and database system.

项目地址:https://gitcode.com/gh_mirrors/ar/arkime
点击查看免费下载

Arkime 的 Kafka 写入插件(位于 capture/plugins/kafka/README.md)将捕获引擎生成的 SPI(Session Packet Index,会话数据索引)流改道发送到 Apache Kafka,而不是直接写入 Elasticsearch。本文以该插件为绝对主线,完整讲解其构建方式、全部配置项、三种消息格式的区别,并结合 capture/plugins/kafka/kafka.c 与 capture/db.c 的源码级实现,说明插件加载、librdkafka 配置、批量发送机制与退出清理的底层原理。读完本文,你将能够独立完成 Arkime 的 Kafka 构建、配置与排障。

插件定位:把 SPI 写入 Kafka,而不是 Elasticsearch

Kafka 写入插件的作用非常明确:Arkime 捕获引擎在会话结束后,会把每条会话的索引数据(SPI)以 JSON 形式批量组装,默认通过 HTTP 发送给 Elasticsearch 的 Bulk API 完成入库;启用该插件后,这段数据流被整体切换为 Kafka 生产者消息,由下游消费者(如自定义处理管道、流式分析平台)接管。

需要特别强调的是,Kafka 插件只接管"会话索引数据(SPI)"的写入通道,并不能替代 Elasticsearch 的全部职责。原文档明确提醒:

Please note that communication to Elasticsearch is still needed, for the stats and other housekeeping tasks.

即:统计(stats)以及各类后台维护任务(housekeeping)仍然需要与 Elasticsearch 保持通信。这意味着部署该插件时,Elasticsearch 集群依然要保留并正常运行,只是会话数据的落库位置从 ES Bulk API 换成了 Kafka topic。

构建:启用 Kafka 支持

Kafka 插件依赖 C 语言客户端库 librdkafka,构建时需要在 easybutton 构建脚本中显式开启对应开关:

./easybutton-build.sh --kafka

从构建脚本 easybutton-build.sh 可以看到,--kafka开关背后做了两类工作:

  • 依赖安装:在 Debian/Ubuntu 上安装librdkafka-dev,在 CentOS/RHEL 系上安装librdkafka-devel(见 easybutton-build.sh),macOS 则通过 brew/port 安装librdkafka(easybutton-build.sh);
  • 编译链接:为configure传入KAFKA_CFLAGS="-I/usr/include/librdkafka/"、KAFKA_LIBS="-lrdkafka",并将--with-kafka置为可用状态;若本机缺少系统 librdkafka,脚本还会自动下载源码并静态编译 librdkafka(easybutton-build.sh),最后单独进入 capture/plugins/kafka 目录完成插件本体的编译(easybutton-build.sh)。

需要说明的是,Kafka 支持默认并未开启,未传--kafka时构建脚本会把--with-kafka置为no(easybutton-build.sh),因此想用该插件必须显式开启此开关。

配置参数全表

Kafka 插件的全部配置项如下表(表格内容完全继承自 capture/plugins/kafka/README.md):

PropertyDetailsExample
kafkaBootstrapServersbootstrap servers, comma separated, to connect to1.2.3.4:9020,5.6.7.8:9020
kafkaTopictopic to send the SPI toarkime-spi
kafkaSSLwhether to enable SSL security protocoltrue
kafkaSSLCALocationpath where the SSL CA is located/path/to/ca.crt
kafkaSSLCertificateLocationpath where the SSL client certificate is located/path/to/client.crt
kafkaSSLKeyLocationpath where the SSL client key is located/path/to/client.key
kafkaSSLKeyPasswordoptional password for the client key
kafkaMsgFormathow to send the SPI data: bulk (default, raw bulk msg), bulk1 (bulk formatted, but just 1 doc), doc (just the doc)bulk

这些参数在 kafka.c 的arkime_plugin_init()中被逐一读取并映射到 librdkafka 配置,下面逐个说明底层行为。

kafkaBootstrapServers:Broker 地址列表

插件通过arkime_config_str_list()读取该参数(kafka.c)。一个容易被忽略的细节是:Arkime 的配置列表使用;作为分隔符,而 librdkafka 的metadata.broker.list要求逗号分隔,因此源码先用arkime_config_str_list()按;拆分,再用g_strjoinv(",", ...)拼回逗号分隔串(kafka.c)后才写入 librdkafka 配置(kafka.c)。所以配置时可按 Arkime 惯例用;分隔多个 broker,插件会自动转换为 librdkafka 需要的格式。

kafkaTopic:SPI 写入的 Topic

默认值为arkime-json(kafka.c)。注意:原文档示例给的是arkime-spi,实际未配置时插件使用arkime-json,生产环境中建议显式配置以保持一致。

kafkaSSL 及证书族:启用 TLS 加密

kafkaSSL是布尔开关(arkime_config_boolean,默认FALSE)。置为true时,插件将 librdkafka 的security.protocol设为SSL(kafka.c),并按需设置:

  • kafkaSSLCALocation→ssl.ca.location(CA 证书路径)
  • kafkaSSLCertificateLocation→ssl.certificate.location(客户端证书路径)
  • kafkaSSLKeyLocation→ssl.key.location(客户端私钥路径)
  • kafkaSSLKeyPassword→ssl.key.password(私钥口令,可选)

这四个 SSL 相关路径参数均为可选:只有配置了才会写入对应 librdkafka 属性(kafka.c)。若rd_kafka_conf_set返回非RD_KAFKA_CONF_OK,插件会直接以LOGEXIT终止(kafka.c)。

进阶:kafka-config 小节直通 librdkafka 属性

除上述固定参数外,kafka.c 还支持一个通用扩展机制:配置文件中[kafka-config]小节里的任意键值对,都会被逐一转发给 librdkafka(kafka.c):

[kafka-config] queue.buffering.max.messages=100000 linger.ms=5

其实现方式是arkime_config_section_keys()枚举小节内所有键,再用rd_kafka_conf_set()逐个写入(kafka.c)。这样无需修改源码即可调优生产者行为(如批量聚合、队列上限、重试策略等),可参考 librdkafka 的配置项文档按需使用。

kafkaMsgFormat:三种 SPI 消息格式

kafkaMsgFormat决定消息的封装形态,默认bulk。结合 db.c 的arkime_db_set_send_bulk2()实现,三种取值对应完全不同的落盘语义:

取值含义底层效果
bulk(默认)原始 Bulk 消息sendBulkHeader=TRUE, indexInDoc=FALSE, maxDocs=0xffff:每条 Kafka 消息内包含{"index":{...}}索引头 + 文档体,一批最多 65535 条文档
bulk1Bulk 格式,但每条仅 1 个文档sendBulkHeader=TRUE, indexInDoc=FALSE, maxDocs=1:保留 Bulk 头、单个文档、每条消息 1 条记录
doc仅文档本身sendBulkHeader=FALSE, indexInDoc=TRUE, maxDocs=1:去掉 Bulk 索引头,把索引名以"index":"..."字段内嵌进文档,每条消息 1 条文档

三种模式的实际分支在 kafka.c,若传入其他值会以LOGEXIT("Unknown config kafkaMsgFormat value ...")直接退出。建议先明确下游消费端期望的格式再选择:

  • 下游按 Elasticsearch Bulk API 语义解析 → 用bulk;
  • 下游需要逐条处理、但想保留 Bulk 头 → 用bulk1;
  • 下游只想消费纯净的会话 JSON(索引信息内嵌)→ 用doc。

配置示例:一份完整的 Kafka 写入配置

结合以上参数,一份典型的 capture 配置(追加到 Arkime 的config.ini中,参考 release/config.ini.sample 的编写风格)如下:

# 开启 Kafka 写入插件 plugins=kafka.so # Kafka 插件参数 kafkaBootstrapServers=1.2.3.4:9020;5.6.7.8:9020 kafkaTopic=arkime-spi kafkaMsgFormat=bulk # 可选:启用 TLS # kafkaSSL=true # kafkaSSLCALocation=/path/to/ca.crt # kafkaSSLCertificateLocation=/path/to/client.crt # kafkaSSLKeyLocation=/path/to/client.key # kafkaSSLKeyPassword=yourpassword # 可选:透传 librdkafka 高级参数 # [kafka-config] # queue.buffering.max.messages=100000

注意:kafkaBootstrapServers中的;分隔会被插件自动转换为逗号再交给 librdkafka;如果直接写逗号分隔也可以,因为arkime_config_str_list支持;或,作为列表分隔,关键在于最终传给 librdkafka 的串必须是逗号分隔(kafka.c)。kafkaSSL未配置时默认关闭,不配置任何 SSL 相关项即可跑明文模式。

源码原理:从会话落库到 Kafka 生产者的完整链路

插件加载与生命周期

arkime_plugin_init()是插件统一入口:先用arkime_plugins_register("kafka", TRUE)注册插件,再通过arkime_plugins_set_cb()注册退出回调kafka_plugin_exit(kafka.c)。随后完成 librdkafka 配置对象创建与上述全部参数注入,最后rd_kafka_new(RD_KAFKA_PRODUCER, ...)创建生产者实例(kafka.c);若创建失败同样LOGEXIT终止。

会话数据如何切到 Kafka:arkime_db_set_send_bulk2 注入点

Arkime 的会话数据库模块在 capture/db.c 维护一个可替换的批量发送函数指针:

LOCAL ArkimeDbSendBulkFunc sendBulkFunc = arkime_db_send_bulk; // 默认发往 ES LOCAL gboolean sendBulkHeader = TRUE; LOCAL gboolean sendIndexInDoc = FALSE; LOCAL uint16_t sendMaxDocs = 0xffff;

arkime_db_set_send_bulk2(func, bulkHeader, indexInDoc, maxDocs)(db.c)就是替换这四个状态量的唯一入口。Kafka 插件在初始化时按kafkaMsgFormat调用它,把sendBulkFunc换成自己的kafka_send_session_bulk,同时设定 Bulk 头、索引内嵌与单批文档上限(kafka.c)。

会话落库的主路径位于 db.c:当缓冲区剩余空间不足或dbInfo[thread].cnt >= sendMaxDocs时触发一次批量发送,调用sendBulkFunc(json, len);若sendBulkHeader为真,则在每条记录前追加{"index":{"_index":"...sessions3-<prefix>", "_id": ...}}行(db.c);若sendIndexInDoc为真,则把索引名写进文档内部的"index"字段(db.c)。这正好从底层解释了上一节三种消息格式的差异来源——Bulk 头与内嵌索引由同一组开关控制,而maxDocs(0xffff/1)直接决定单条 Kafka 消息携带的文档数量。批量缓冲区大小受全局dbBulkSize约束(默认 1000000,范围 500000~15000000,见 config.c)。

生产者发送与队列背压处理

核心发送函数kafka_send_session_bulk(kafka.c)使用rd_kafka_producev将 JSON 消息压入生产者队列,并做两层处理:

  • 投递成功:立即以非阻塞方式rd_kafka_poll(rk, 0)驱动回调;
  • RD_KAFKA_RESP_ERR__QUEUE_FULL(队列满):说明 librdkafka 内部队列已满(受queue.buffering.max.messages限制),插件会rd_kafka_poll(rk, 100)阻塞等待最多 100ms 后重试,最多重试 5 次(kafka.c);其他错误则直接放弃不再重试。

每次消息的 JSON 缓冲区作为_private(V_OPAQUE)随消息传入,投递回调kafka_msg_delivered_bulk_cb在确认送达后调用arkime_http_free_buffer(json)释放(kafka.c);若最终未能发送,也在发送函数尾部释放(kafka.c),避免内存泄漏。该回调还会在config.debug开启时打印投递耗时、字节数、offset、partition 与 broker 等诊断信息,debug > 3时甚至打印完整 payload,是排查"消息是否真正送达"的第一现场。

退出清理:Flush 与未投递统计

kafka_plugin_exit()(kafka.c)在 Arkime 退出时被调用:先rd_kafka_flush(rk, 10*1000)最多等待 10 秒冲刷剩余消息;若rd_kafka_outq_len(rk) > 0,说明仍有消息未投递,会打印未送达数量告警;随后rd_kafka_destroy(rk)销毁生产者实例。这意味着正常情况下退出前队列会尽量排空,但若 10 秒内未能全部送达,日志中会出现明确的未投递计数,可作为数据完整性评估依据。

验证与排障要点

  • 确认插件已加载:启动 capture 时日志出现Loading Kafka plugin与Kafka plugin loaded(kafka.c);
  • 确认配置合法:任一rd_kafka_conf_set失败都会以LOGEXIT立即退出并打印错误串(如Error configuring kafka:metadata.broker.list, error = ...),这是定位参数拼写错误的最直接手段(kafka.c);
  • 确认消息真正送达:开启 debug(debug=1)后可看到Message delivered in ... ms (N bytes, offset ..., partition ..., broker ...)投递报告;debug > 3还会输出 payload(kafka.c);
  • 观察队列背压:若日志反复出现Failed to produce to topic ...: Local: Queue full,说明生产者吞吐跟不上 broker,可调大[kafka-config]中的queue.buffering.max.messages或增加分区/消费者;
  • 退出时留意未投递告警:N message(s) were not delivered意味着 flush 窗口(10 秒)内未全部送出,需要结合 broker 状态判断是否丢数(kafka.c)。

小结

Kafka 写入插件为 Arkime 提供了"SPI 走 Kafka、统计与维护仍走 Elasticsearch"的混合架构:只需./easybutton-build.sh --kafka重新构建,并在配置中指定kafkaBootstrapServers、kafkaTopic与kafkaMsgFormat,即可把会话索引数据流切换到 Kafka。理解bulk/bulk1/doc三种格式与arkime_db_set_send_bulk2()的对应关系(db.c),并结合 kafka.c 的队列背压、投递回调与退出 flush 机制进行调优与排障,就能在生产环境中稳定地让 Arkime 与会话数据的 Kafka 化无缝衔接。

  • 网络安全
  • 网络
  • 后端
  • 数据可视化

【免费下载链接】arkime

Arkime is an open source, large scale, full packet capturing, indexing, and database system.

项目地址:https://gitcode.com/gh_mirrors/ar/arkime
点击查看免费下载
上一篇:Make Me a Hanzi:开源汉字学习与数据资源终极指南
下一篇:自然语言处理在文化遗产数字化保护中的10大应用场景:完整指南

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询