kafka-ui 如何注册 S3 Sink Connector 并将主题消息写入 S3 存储桶? kafka-ui 如何注册 S3 Sink Connector 并将主题消息写入 S3 存储桶【免费下载链接】kafka-uiOpen-Source Web UI for Apache Kafka Management项目地址: https://gitcode.com/GitHub_Trending/ka/kafka-ui这篇文档针对一个具体任务把 Kafka 主题中的消息通过 S3 Sink Connector 落盘到 S3 存储桶并让 kafka-ui 能管理这个 Connector。操作依据是仓库documentation/compose/下自带的 e2e 环境文件它同时给出了 Kafka Connect 镜像含 S3 connector 插件、Connector 注册脚本和 connector 配置 JSON。完成后的结果是github-issues、github-pull_requests、github-commits三个主题中的消息被写入配置中指定的 S3 桶文档示例为kafka-ui-s3-sink-connectorregioneu-central-1。环境组成e2e-tests.yaml 里有哪些组件整个链路由 e2e-tests.yaml 编排DOCKER_COMPOSE.md 对它的描述是包含不同 connectorsgithub-source、s3、sink-activities、source-activities和 Ksql 功能的配置。与当前任务相关的服务有kafka-connect0基于 Dockerfile 构建镜像基础是confluentinc/cp-kafka-connect:6.0.1构建时通过confluent-hub install --no-prompt confluentinc/kafka-connect-s3:latest安装 S3 connector 插件同时安装了 jdbc 和 github 插件对宿主机暴露8083:8083。kafka0confluentinc/cp-kafka:7.2.1KRaft 模式无 Zookeeper。schemaregistry0confluentinc/cp-schema-registry:7.2.1。kafka-uiprovectuslabs/kafka-ui:latest其中两个环境变量把 Connect 接入 UIKAFKA_CLUSTERS_0_KAFKACONNECT_0_NAME: first KAFKA_CLUSTERS_0_KAFKACONNECT_0_ADDRESS: http://kafka-connect0:8083kafka-ui 通过这两个变量KAFKACONNECT_*_NAME/KAFKACONNECT_*_ADDRESS连接 Kafka Connect之后可以在 UI 中从 connectors 视图跳转到对应主题README 中展示了 Connector、Topic、Consumer 之间互相跳转的导航但 connector 本身的注册动作是通过 Connect 的 REST API 完成的下面一节说明。配置 S3 Sink ConnectorConnector 配置在 s3-sink.json仓库中的完整内容为{ name: s3-sink, config: { connector.class: io.confluent.connect.s3.S3SinkConnector, topics: github-issues, github-pull_requests, github-commits, tasks.max: 1, s3.region: eu-central-1, s3.bucket.name: kafka-ui-s3-sink-connector, s3.part.size: 5242880, flush.size: 3, storage.class: io.confluent.connect.s3.storage.S3Storage, format.class: io.confluent.connect.s3.format.json.JsonFormat, schema.generator.class: io.confluent.connect.storage.hive.schema.DefaultSchemaGenerator, partitioner.class: io.confluent.connect.storage.partitioner.DefaultPartitioner, schema.compatibility: BACKWARD } }其中影响你自己环境的取值s3.bucket.name与s3.region文档示例值分别为kafka-ui-s3-sink-connector和eu-central-1实际使用时替换成你已创建的桶名和所在 region。topics三个由 github-source connector 写入的主题见下文如果消息来源不同替换为你自己的主题列表。其余类名storage.class、format.class、schema.generator.class、partitioner.class是文档给出的示例取值仓库文档未逐一解释其含义照抄即可跑通该示例。AWS 凭据不在 JSON 里而是在e2e-tests.yaml的kafka-connect0环境变量的最后两行文档中以注释形式给出、值为空# AWS_ACCESS_KEY_ID: # AWS_SECRET_ACCESS_KEY: 需要取消注释并填入有该桶写权限的凭据后再启动 compose 环境否则 connector 无法访问 S3。注册 ConnectorPOST 到 Connect 的 /connectors 端点注册动作就是把上述 JSON POST 到 Kafka Connect REST 端点。仓库自带两种方式方式一单条命令只注册 s3-sink本文主路径e2e-tests.yaml中kafka-connect0映射了8083:8083在宿主机上执行curl -X POST -H Content-Type: application/json \ -d documentation/compose/connectors/s3-sink.json \ http://localhost:8083/connectors如果s3-sink.json在你本地是其他路径替换-d 后的文件路径即可。方式二compose 的 create-connectors 服务批量注册可选e2e-tests.yaml中定义了create-connectors服务镜像ellerbrock/alpine-bash-curl-ssl把./connectors目录挂载为/connectors执行bash -c /connectors/start.sh。start.sh 的逻辑是while [[ $(curl -s -o /dev/null -w %{http_code} kafka-connect0:8083) ! 200 ]] do sleep 5 done echo \n --------------Creating connectors... for filename in /connectors/*.json; do curl -X POST -H Content-Type: application/json -d $filename http://kafka-connect0:8083/connectors done它先轮询直到kafka-connect0:8083返回 HTTP 200再把目录下所有.json逐一 POST 到/connectors。注意这个脚本会把同目录下的github-source.json、s3-sink.json、sink-activities.json、source-activities.json全部注册包括 PostgreSQL 相关 connector因此它只适合整套 e2e 环境一起跑只想注册 S3 sink 时请用方式一。消息从哪里来github-source 写入三个主题s3-sink消费的三个主题由同目录的 github-source.json 产生github.resources:issues,commits,pull_requeststopic.name.pattern:github-${resourceName}即按该模式生成github-issues、github-commits、github-pull_requests三个主题topic.name.pattern中的${resourceName}是 Confluent GitHub connector 自身的模板变量由该 connector 按其 resources 配置展开不需要手动替换。S3 sink 的topics字段与这三个主题一一对应。该 connector 还需要github.access.token文档示例中为空实际运行需填入有效的 GitHub token。若你只想验证 S3 写入链路而不需要 GitHub 数据源可以用自己的生产者往同名或改配置后对应的主题写消息。结果验证与边界文档中给出的可核对状态是 Connect 的可用性检查start.sh把curl到kafka-connect0:8083返回200作为创建 connector 前的就绪条件kafka-ui 本身的健康检查端点是/actuator/healthREADME 中说明 liveliness/readiness 位于/actuator/health。注册成功后connector 会按配置把三个主题的消息写入对应 region 下的s3.bucket.name桶仓库文档没有给出 S3 侧的输出示例或对象路径样例确认数据到达需要在 S3 控制台或客户端中按你配置的 region/桶查看具体路径结构由 connector 的 partitioner/format 配置决定文档未展开。需要留意的限制AWS 凭据在e2e-tests.yaml中默认是注释掉的空值不填写则无法写入 S3。kafka-connect0插件路径CONNECT_PLUGIN_PATH: /usr/share/java,/usr/share/confluent-hub-components与 Dockerfile 中confluent-hub的安装位置对应S3 插件必须通过该 Dockerfile 构建的镜像提供。该 compose 文件定位为 e2e/演示环境KRaft 单节点、无认证不是生产部署方案。【免费下载链接】kafka-uiOpen-Source Web UI for Apache Kafka Management项目地址: https://gitcode.com/GitHub_Trending/ka/kafka-ui创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考