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

如何在GCP Dataproc上配置PySpark Streaming对接Cloudkarafka

依赖服务说明
  • GCP Dataproc & DataprocHub
  • Cloudkarafka:可快速搭建Kafka服务的轻量化托管产品
详细操作步骤
  1. GCP DataprocHub默认预装Spark 2.4.8版本。
  2. 创建指定版本的Dataproc集群:
# gcloud shell执行
gcloud dataproc clusters create dataproc-spark312  --image-version=2.0-ubuntu18  --region=us-central1 --single-node
  1. 导出集群配置并存入GS存储桶:
gcloud dataproc clusters export dataproc-spark312 --destination dataproc-spark312.yaml --region us-central1
gsutil cp dataproc-spark312.yaml gs://gcp-learn-lib/
  1. 创建环境变量配置文件:
DATAPROC_CONFIGS=gs://gcp-learn-lib/dataproc-spark312.yaml
NOTEBOOKS_LOCATION=gs://gcp-learn-notebooks/notebooks
DATAPROC_LOCATIONS_LIST=a,b,c

保存为dataproc-hub-config.env后上传至GS存储桶。
5. 创建Dataproc Hub并关联上述集群,在「自定义环境配置」区块添加配置:key为container-env-file,value为gs://gcp-learn-lib/dataproc-hub-config.env。
6. 完成Dataproc Hub创建流程。
7. 点击Jupyter访问链接。
8. 页面将展示已关联集群,选中集群并选择与步骤2一致的区域。
9. 进入Dataproc -> 集群 -> 点击名称以hub-开头的集群 -> 虚拟机实例 -> SSH,下载Kafka依赖Jar包到指定目录:

cd /usr/lib/spark/jars/
wget https://repo1.maven.org/maven2/org/apache/commons/commons-pool2/2.6.2/commons-pool2-2.6.2.jar
wget https://repo1.maven.org/maven2/org/apache/kafka/kafka-clients/2.6.0/kafka-clients-2.6.0.jar
wget https://repo1.maven.org/maven2/org/apache/spark/spark-sql-kafka-0-10_2.12/3.1.2/spark-sql-kafka-0-10_2.12-3.1.2.jar
wget https://repo1.maven.org/maven2/org/apache/spark/spark-streaming-kafka-0-10_2.12/3.1.2/spark-streaming-kafka-0-10_2.12-3.1.2.jar
wget https://repo1.maven.org/maven2/org/apache/spark/spark-streaming-kafka-0-10-assembly_2.12/3.1.2/spark-streaming-kafka-0-10-assembly_2.12-3.1.2.jar
  1. 准备JAAS配置文件与CA证书文件:
vi /tmp/cloudkarafka_gcp_oct2021.jaas
vi /tmp/cloudkarafka_gcp_oct2021.ca
  1. 向Spark默认配置文件末尾追加如下内容:
# 先打开配置文件:sudo vi /etc/spark/conf/spark-defaults.conf
spark.executor.extraJavaOptions=-Djava.security.auth.login.config=/tmp/cloudkarafka_gcp_oct2021.jaas -Dsasl.jaas.config=/tmp/cloudkarafka_gcp_oct2021.jaas -Dssl.ca.location=/tmp/cloudkarafka_gcp_oct2021.ca
spark.driver.extraJavaOptions=-Djava.security.auth.login.config=/tmp/cloudkarafka_gcp_oct2021.jaas
  1. 返回JupyterHub,重启内核后执行如下测试代码验证连通性:
spark.readStream.format("kafka") \
.option("kafka.bootstrap.servers", "server1:9094,serer2:9094,server3:9094") \
.option("subscribe", "foo-default") \
.option("kafka.security.protocol", "SASL_SSL") \
.option("kafka.sasl.mechanism", "SCRAM-SHA-256") \
.load() \
.selectExpr("CAST(key AS STRING)", "CAST(value AS STRING)") \
.writeStream.format("kafka") \
.option("kafka.bootstrap.servers", "server1:9094,serer2:9094,server3:9094") \
.option("topic", "foo-test") \
.option("kafka.security.protocol", "SASL_SSL") \
.option("kafka.sasl.mechanism", "SCRAM-SHA-256") \
.option("checkpointLocation", "/tmp/stream/kafkatest") \
.start()
  1. 配置完成即可正常使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 10:09:03