Apache APISIX kafka-proxy 插件:为 Kafka Upstream 配置 SASL/PLAIN 认证的完整指南
2026/9/15 0:34:41 网站建设 项目流程

Apache APISIX kafka-proxy 插件:为 Kafka Upstream 配置 SASL/PLAIN 认证的完整指南

【免费下载链接】apisixThe Cloud-Native API Gateway项目地址: https://gitcode.com/GitHub_Trending/ap/apisix

Apache APISIX 的kafka-proxy插件用于为scheme 为kafka的 Upstream配置高级参数,核心能力是为 Kafka 客户端注入SASL/PLAIN 用户名密码认证。本文将基于官方文档并结合仓库源码,完整讲解该插件的配置属性、加密存储机制、底层调用链与实战用法,帮助你在通过 WebSocket 桥接 Kafka 时安全地完成认证配置。

插件概述

kafka-proxy插件本身不负责 Kafka 协议代理,而是充当"配置注入器":它把 SASL 认证信息写入当前请求的上下文(ctx),随后 APISIX 的 Kafka 桥接模块会读取这些上下文变量,为底层 Kafka 客户端组装sasl_config。从源码结构看,该插件实际配合 apisix/pubsub/kafka.lua 这一 Pub/Sub 桥接实现一起工作,后者才是真正接管 Kafka 请求、执行list_offsetfetch等命令的模块。

插件的注册信息位于 apisix/plugins/kafka-proxy.lua:

local _M = { version = 0.1, priority = 508, name = "kafka-proxy", schema = schema, }

其优先级为508,在请求处理早期(access阶段)执行,确保认证上下文在 Kafka 桥接逻辑触发前就绪。

属性(Attributes)

kafka-proxy的属性非常精简,核心是一个嵌套的sasl对象:

NameTypeRequiredDefaultValid valuesDescription
saslobjectoptional{"username": "user", "password": "pwd"}SASL/PLAIN 认证配置;当该配置存在时即开启 SASL 认证;该对象包含 username 与 password 两个参数,且二者都必须配置。
sasl.usernamestringrequiredSASL/PLAIN 认证用户名
sasl.passwordstringrequiredSASL/PLAIN 认证密码

对应的 JSON Schema 定义(见 apisix/plugins/kafka-proxy.lua):

local schema = { type = "object", properties = { sasl = { type = "object", properties = { username = { type = "string" }, password = { type = "string" }, }, required = {"username", "password"}, }, }, encrypt_fields = {"sasl.password"}, }

关键要点:

  • sasl为可选对象,一旦配置则 username 和 password 均为必填,缺一不可;
  • username 和 password 的类型都必须是string,传入数字等非字符串类型会被 Schema 校验拒绝;
  • 当前 SASL 认证仅支持 PLAIN 模式(即用户名密码登录方式),不支持 SCRAM、GSSAPI 等其他机制。

这些约束在测试用例 t/plugin/kafka-proxy.t 的TEST 1: sanity中有直接验证:空配置和完整配置通过校验,而缺少password或密码为数字时会分别报出property "password" is requiredwrong type: expected string, got number的校验错误。

密码加密存储(encrypt_fields)

Schema 中声明了encrypt_fields = {"sasl.password"},这意味着sasl.password会以加密形式存储在 etcd中,而不是明文落盘。该机制属于 APISIX 的数据加密能力,详细设计可参考 插件开发指南 中关于 encrypted storage fields 的章节。

测试 t/plugin/kafka-proxy.t 的TEST 2: data encryption for sasl.password完整验证了这一行为:启用data_encryption后,通过 Admin API 写入路由时密码为明文(admin-secret),查询 Admin API 时返回解密后的明文,而直读 etcd 中/routes/1的存储值时,sasl.password已变成密文(如y4Z3aqo51xrt3f9UziNUrg==)。

# 启用字段加密所需的最小配置(对应测试中的 yaml_config) apisix: data_encryption: enable_encrypt_fields: true keyring: - edd1c9f0985e76a2

工作原理:从插件到 Kafka 客户端的调用链

kafka-proxy之所以能生效,依赖 APISIX 对kafkascheme 上游的特殊分流。整个链路如下:

1. access 阶段注入认证上下文

插件的核心逻辑非常简洁(apisix/plugins/kafka-proxy.lua):

function _M.access(conf, ctx) if conf.sasl then ctx.kafka_consumer_enable_sasl = true ctx.kafka_consumer_sasl_username = conf.sasl.username ctx.kafka_consumer_sasl_password = conf.sasl.password end end

sasl配置存在时,插件把三个变量写入请求上下文:

  • kafka_consumer_enable_sasl:是否启用 SASL 的标志位;
  • kafka_consumer_sasl_username:SASL 用户名;
  • kafka_consumer_sasl_password:SASL 密码。

2. 按 scheme 分流到 Kafka 桥接模块

在 APISIX 的请求入口 apisix/init.lua 中,当匹配到的 Upstream 的scheme == "kafka"时,会跳过负载均衡(注释明确说明 "load balancer is not required by kafka upstream"),直接交由pubsub_kafka.access(api_ctx)接管:

-- load balancer is not required by kafka upstream, so the upstream -- node selection process is intercepted and left to kafka to -- handle on its own if api_ctx.matched_upstream and api_ctx.matched_upstream.scheme == "kafka" then return pubsub_kafka.access(api_ctx) end

3. 桥接模块组装 Kafka 客户端配置

在 apisix/pubsub/kafka.lua 中,桥接模块遍历 Upstream 节点构造 broker 列表,并读取请求上下文中由kafka-proxy插件写入的 SASL 信息,为每个 broker 填充sasl_config

local up_nodes = api_ctx.matched_upstream.nodes local broker_list = {} for i, node in ipairs(up_nodes) do broker_list[i] = { host = node.host, port = node.port, } if api_ctx.kafka_consumer_enable_sasl then broker_list[i].sasl_config = { mechanism = "PLAIN", user = api_ctx.kafka_consumer_sasl_username, password = api_ctx.kafka_consumer_sasl_password, } end end local client_config = {refresh_interval = 30 * 60 * 1000} if api_ctx.matched_upstream.tls then client_config.ssl = true client_config.ssl_verify = api_ctx.matched_upstream.tls.verify end local consumer = bconsumer:new(broker_list, client_config)

值得注意的细节:

  • SASL 机制被硬编码为"PLAIN",与文档中"仅支持 PLAIN 模式"的描述一致;
  • 如果 Upstream 配置了tls字段,桥接模块会同时开启sslssl_verify,因此kafka-proxy可与 Kafka 的 TLS 传输叠加使用;
  • 底层客户端是resty.kafka.basic-consumer,broker 元数据刷新间隔为 30 分钟。

4. 注册 Kafka 命令回调并进入事件循环

桥接模块通过core.pubsub(见 apisix/core/pubsub.lua)注册cmd_kafka_list_offsetcmd_kafka_fetch两个命令回调,分别处理 Kafka 的"查询分区 offset"与"拉取消息"操作,随后调用pubsub:wait()进入事件循环等待 WebSocket 客户端发送命令。kafka之所以可以作为 Upstream 的 scheme,是因为它在 apisix/schema_def.lua 的scheme枚举中被显式列出(与grpc/grpcs/http/https/tcp/tls/udp并列,描述为 "For specific protocols, it can be kafka")。

使用示例

当 Upstream 使用kafkascheme 时,可以通过kafka-proxy插件为其附加 Kafka 认证配置。下面创建一个名为r1的路由,将所有发往/kafka的请求代理到三个 Kafka broker,并开启 SASL/PLAIN 认证:

curl -X PUT 'http://127.0.0.1:9180/apisix/admin/routes/r1' \ -H 'X-API-KEY: <api-key>' \ -H 'Content-Type: application/json' \ -d '{ "uri": "/kafka", "plugins": { "kafka-proxy": { "sasl": { "username": "user", "password": "pwd" } } }, "upstream": { "nodes": { "kafka-server1:9092": 1, "kafka-server2:9092": 1, "kafka-server3:9092": 1 }, "type": "none", "scheme": "kafka" } }'

配置说明:

  • plugins.kafka-proxy.sasl:开启 SASL/PLAIN 认证,usernamepassword必须同时给出;
  • upstream.scheme:必须为kafka,这是 APISIX 分流到 Kafka 桥接模块的判定依据;
  • upstream.type:官方示例使用none,因为 Kafka 场景下负载均衡被桥接模块接管(见上文 apisix/init.lua 的分流逻辑);
  • upstream.nodes:以host:port形式列出 Kafka broker 地址,权重值(示例中为 1)在此场景下主要用于节点列举。

创建完成后,即可通过WebSocket连接/kafka端点进行测试。WebSocket 是 APISIX 桥接 Kafka 的客户端通道:客户端在 WebSocket 上发送cmd_kafka_list_offsetcmd_kafka_fetch等二进制命令,由 apisix/pubsub/kafka.lua 中的回调执行真实 Kafka 操作并回传结果。由于桥接模块支持sslssl_verify(来自 Upstream 的tls配置),SASL 认证可与 TLS 加密传输同时生效,适合生产环境中的安全要求。

删除插件

移除kafka-proxy插件时,只需从路由的插件配置中删除对应的 JSON 片段即可。例如再次调用 Admin API,将路由中plugins下的kafka-proxy对象移除:

curl -X PUT 'http://127.0.0.1:9180/apisix/admin/routes/r1' \ -H 'X-API-KEY: <api-key>' \ -H 'Content-Type: application/json' \ -d '{ "uri": "/kafka", "plugins": {}, "upstream": { "nodes": { "kafka-server1:9092": 1 }, "type": "none", "scheme": "kafka" } }'

APISIX 会自动热加载新的配置,无需重启即可生效。删除后,请求上下文中的kafka_consumer_enable_sasl将不再被置位,Kafka 客户端将按无认证模式连接 broker。

使用限制与注意事项

  1. 仅支持 SASL/PLAIN:当前实现将mechanism硬编码为"PLAIN",若 Kafka 集群配置了 SCRAM 等其他 SASL 机制,该插件无法适配;
  2. 认证与加密传输独立:SASL 认证解决身份问题,传输加密需另行在 Upstream 中配置tls字段;
  3. 敏感信息默认加密存储:由于sasl.password声明了encrypt_fields,建议在配置中启用data_encryption,避免密码明文写入 etcd;
  4. Schema 校验严格sasl一旦出现,usernamepassword必须为字符串且必填,任何缺失或类型错误都会导致配置被拒绝(可参考 t/plugin/kafka-proxy.t 中的校验断言);
  5. 配套组件:该插件只有在 Upstreamschemekafka时才有实际意义,普通 HTTP/TCP 上游无需也不应配置此插件。

延伸阅读

  • 插件源码:apisix/plugins/kafka-proxy.lua
  • Kafka 桥接实现:apisix/pubsub/kafka.lua
  • Pub/Sub 基础模块:apisix/core/pubsub.lua
  • Kafka scheme 分流逻辑:apisix/init.lua
  • Upstream scheme 枚举定义:apisix/schema_def.lua
  • 加密存储字段机制:插件开发指南
  • 测试用例:t/plugin/kafka-proxy.t

【免费下载链接】apisixThe Cloud-Native API Gateway项目地址: https://gitcode.com/GitHub_Trending/ap/apisix

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

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

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

立即咨询