完全指南:配置、构建与原理)
网络安全网络后端数据可视化【免费下载链接】arkimeArkime 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将捕获引擎生成的 SPISession Packet Index会话数据索引流改道发送到 Apache Kafka而不是直接写入 Elasticsearch。本文以该插件为绝对主线完整讲解其构建方式、全部配置项、三种消息格式的区别并结合 capture/plugins/kafka/kafka.c 与 capture/db.c 的源码级实现说明插件加载、librdkafka 配置、批量发送机制与退出清理的底层原理。读完本文你将能够独立完成 Arkime 的 Kafka 构建、配置与排障。插件定位把 SPI 写入 Kafka而不是 ElasticsearchKafka 写入插件的作用非常明确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.shmacOS 则通过 brew/port 安装librdkafkaeasybutton-build.sh编译链接为configure传入KAFKA_CFLAGS-I/usr/include/librdkafka/、KAFKA_LIBS-lrdkafka并将--with-kafka置为可用状态若本机缺少系统 librdkafka脚本还会自动下载源码并静态编译 librdkafkaeasybutton-build.sh最后单独进入 capture/plugins/kafka 目录完成插件本体的编译easybutton-build.sh。需要说明的是Kafka 支持默认并未开启未传--kafka时构建脚本会把--with-kafka置为noeasybutton-build.sh因此想用该插件必须显式开启此开关。配置参数全表Kafka 插件的全部配置项如下表表格内容完全继承自 capture/plugins/kafka/README.mdPropertyDetailsExamplekafkaBootstrapServersbootstrap servers, comma separated, to connect to1.2.3.4:9020,5.6.7.8:9020kafkaTopictopic to send the SPI toarkime-spikafkaSSLwhether to enable SSL security protocoltruekafkaSSLCALocationpath where the SSL CA is located/path/to/ca.crtkafkaSSLCertificateLocationpath where the SSL client certificate is located/path/to/client.crtkafkaSSLKeyLocationpath where the SSL client key is located/path/to/client.keykafkaSSLKeyPasswordoptional password for the client keykafkaMsgFormathow 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 配置下面逐个说明底层行为。kafkaBootstrapServersBroker 地址列表插件通过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 需要的格式。kafkaTopicSPI 写入的 Topic默认值为arkime-jsonkafka.c。注意原文档示例给的是arkime-spi实际未配置时插件使用arkime-json生产环境中建议显式配置以保持一致。kafkaSSL 及证书族启用 TLS 加密kafkaSSL是布尔开关arkime_config_boolean默认FALSE。置为true时插件将 librdkafka 的security.protocol设为SSLkafka.c并按需设置kafkaSSLCALocation→ssl.ca.locationCA 证书路径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]小节里的任意键值对都会被逐一转发给 librdkafkakafka.c[kafka-config] queue.buffering.max.messages100000 linger.ms5其实现方式是arkime_config_section_keys()枚举小节内所有键再用rd_kafka_conf_set()逐个写入kafka.c。这样无需修改源码即可调优生产者行为如批量聚合、队列上限、重试策略等可参考 librdkafka 的配置项文档按需使用。kafkaMsgFormat三种 SPI 消息格式kafkaMsgFormat决定消息的封装形态默认bulk。结合 db.c 的arkime_db_set_send_bulk2()实现三种取值对应完全不同的落盘语义取值含义底层效果bulk默认原始 Bulk 消息sendBulkHeaderTRUE, indexInDocFALSE, maxDocs0xffff每条 Kafka 消息内包含{index:{...}}索引头 文档体一批最多 65535 条文档bulk1Bulk 格式但每条仅 1 个文档sendBulkHeaderTRUE, indexInDocFALSE, maxDocs1保留 Bulk 头、单个文档、每条消息 1 条记录doc仅文档本身sendBulkHeaderFALSE, indexInDocTRUE, maxDocs1去掉 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 写入插件 pluginskafka.so # Kafka 插件参数 kafkaBootstrapServers1.2.3.4:9020;5.6.7.8:9020 kafkaTopicarkime-spi kafkaMsgFormatbulk # 可选启用 TLS # kafkaSSLtrue # kafkaSSLCALocation/path/to/ca.crt # kafkaSSLCertificateLocation/path/to/client.crt # kafkaSSLKeyLocation/path/to/client.key # kafkaSSLKeyPasswordyourpassword # 可选透传 librdkafka 高级参数 # [kafka-config] # queue.buffering.max.messages100000注意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_exitkafka.c。随后完成 librdkafka 配置对象创建与上述全部参数注入最后rd_kafka_new(RD_KAFKA_PRODUCER, ...)创建生产者实例kafka.c若创建失败同样LOGEXIT终止。会话数据如何切到 Kafkaarkime_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 头与内嵌索引由同一组开关控制而maxDocs0xffff/1直接决定单条 Kafka 消息携带的文档数量。批量缓冲区大小受全局dbBulkSize约束默认 1000000范围 500000~15000000见 config.c。生产者发送与队列背压处理核心发送函数kafka_send_session_bulkkafka.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 缓冲区作为_privateV_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 loadedkafka.c确认配置合法任一rd_kafka_conf_set失败都会以LOGEXIT立即退出并打印错误串如Error configuring kafka:metadata.broker.list, error ...这是定位参数拼写错误的最直接手段kafka.c确认消息真正送达开启 debugdebug1后可看到Message delivered in ... ms (N bytes, offset ..., partition ..., broker ...)投递报告debug 3还会输出 payloadkafka.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 化无缝衔接。赞分享网络安全网络后端数据可视化【免费下载链接】arkimeArkime is an open source, large scale, full packet capturing, indexing, and database system.项目地址https://gitcode.com/gh_mirrors/ar/arkime点击查看免费下载相关推荐Apache Druid Kafka Simple Consumerkafka-0.8-v2Firehose 接入指南配置、原理与生产建议Apache Druid Kafka Simple Consumerkafka 0.8 v2Firehose 接入指南配置、原理与生产建议 本指南围绕 A数据库数据分析OLAP大数据实时分析数据仓库后端Arkime capture 接入 Snort DAQreader-daq 插件的构建、配置与源码原理Arkime capture 接入 Snort DAQreader daq 插件的构建、配置与源码原理 导读 本文围绕 capture/plugins/daq网络安全网络后端数据可视化Apache APISIX kafka-proxy 插件为 Kafka Upstream 配置 SASL/PLAIN 认证的完整指南Apache APISIX kafka proxy 插件为 Kafka Upstream 配置 SASL/PLAIN 认证的完整指南 Apache APISI后端微服务云原生上一篇Make Me a Hanzi开源汉字学习与数据资源终极指南下一篇自然语言处理在文化遗产数字化保护中的10大应用场景完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考