基于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
- 在GCP控制台创建ServiceAccount,赋予
roles/pubsub.subscriber权限到目标Pub/Sub订阅(需先为Pub/Sub主题创建订阅) - 下载ServiceAccount的JSON密钥文件,命名为
key.json - 在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
相关产品推荐
相关产品推荐

