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

基于Strimzi Operator的Kafka Connect对接GCP PubSub配置求助

基于Strimzi搭建GCP Pub/Sub到Kafka的Connector

前置要求

  • 已在GKE上运行Strimzi v0.38 + Kafka 3.5.1集群
  • 拥有GCP权限:能创建ServiceAccount、配置Pub/Sub权限、下载密钥文件
  • 已安装kubectl并能访问目标GKE集群

步骤1:部署KafkaConnect集群(若未部署)

Strimzi通过KafkaConnect资源管理Connector运行环境,需先部署该集群并加载Pub/Sub连接器依赖。以下是示例YAML:

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaConnect
metadata:
  name: pubsub-connect-cluster
  annotations:
    strimzi.io/use-connector-resources: "true"
spec:
  version: 3.5.1
  replicas: 1
  bootstrapServers: <你的Kafka集群BOOTSTRAP地址,比如my-cluster-kafka-bootstrap:9092>
  image: confluentinc/cp-kafka-connect:7.4.0  # 该镜像内置Pub/Sub连接器,版本与Kafka 3.5兼容
  config:
    group.id: pubsub-connect-group
    offset.storage.topic: connect-offsets
    config.storage.topic: connect-configs
    status.storage.topic: connect-status
    key.converter: org.apache.kafka.connect.json.JsonConverter
    value.converter: org.apache.kafka.connect.json.JsonConverter
    key.converter.schemas.enable: false
    value.converter.schemas.enable: false
  externalConfiguration:
    volumes:
      - name: gcp-key
        secret:
          secretName: gcp-pubsub-key

替换<你的Kafka集群BOOTSTRAP地址>为实际地址,可通过kubectl get kafka <你的集群名> -o jsonpath='{.status.listeners[?(@.type=="internal")].bootstrapServers}'获取


步骤2:创建GCP权限Secret

  1. 在GCP控制台创建ServiceAccount,赋予roles/pubsub.subscriber权限到目标Pub/Sub订阅(需先为Pub/Sub主题创建订阅)
  2. 下载ServiceAccount的JSON密钥文件,命名为key.json
  3. 在GKE集群中创建Secret:
kubectl create secret generic gcp-pubsub-key --from-file=key.json=./key.json

步骤3:部署Pub/Sub Source Connector

创建KafkaConnector资源YAML,这是拉取Pub/Sub消息到Kafka的核心配置:

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaConnector
metadata:
  name: pubsub-source-connector
  labels:
    strimzi.io/cluster: pubsub-connect-cluster  # 关联步骤1的KafkaConnect集群名
spec:
  class: com.google.pubsub.kafka.source.PubSubSourceConnector
  tasksMax: 2
  config:
    # GCP基础配置
    pubsub.project.id: <你的GCP项目ID>
    pubsub.subscription.id: <你的Pub/Sub订阅ID>
    gcp.credentials.file.path: /opt/kafka/external-configuration/gcp-key/key.json  # 对应步骤1中Secret的挂载路径
    # Kafka目标配置
    kafka.topic: <要写入的Kafka主题名>
    # 消息拉取配置(可选)
    pubsub.max.batch.size: 100  # 每次拉取的消息数量
    pubsub.poll.timeout: 10000  # 拉取超时时间(毫秒)
    # 偏移量管理(可选)
    pubsub.offset.commit.interval.ms: 5000  # 偏移量提交间隔

替换所有<>中的占位符为实际值


步骤4:验证Connector状态

部署后检查状态:

kubectl get kafkaconnectors pubsub-source-connector

若状态为Ready,则说明Connector已正常运行。可通过查看KafkaConnect Pod日志确认消息拉取情况:

kubectl logs -f <KafkaConnect Pod名>

内容的提问来源于stack exchange,提问作者Melvin

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 02:46:12