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

如何通过Python SDK让SSL证书在Dataflow工作节点可用?

让Dataflow工作节点获取Kafka SSL证书的解决方案

方案一:利用Dataflow文件Staging功能自动同步证书

Dataflow支持通过--files_to_stage参数,将指定的GCS文件自动复制到所有工作节点的当前工作目录,这是最简洁的实现方式。

  1. 提交任务时指定要同步的SSL文件
    在启动Dataflow任务的命令中添加--files_to_stage参数,指向GCS上的证书文件:

    python your_pipeline.py \
      --runner=DataflowRunner \
      --project=你的GCP项目ID \
      --region=你的区域 \
      --temp_location=gs://你的桶/temp \
      --files_to_stage=gs://你的桶/myp12file.p12
    
  2. 修改Kafka消费者配置中的文件路径
    证书会被同步到工作节点的当前工作目录,直接使用文件名即可:

    msg_kv_bytes = pipeline | "read file" >> ReadFromKafka(
         consumer_config={
          "bootstrap.servers": "IP:PORT",
          "group.id": "xyz",
          "default.api.timeout.ms" : "300000",
          "security.protocol": "SSL",
            "ssl.truststore.location": "myp12file.p12",
            "ssl.truststore.password": "***",
            "ssl.truststore.type": "PKCS12",
            "ssl.keystore.type": "PKCS12",
            "ssl.keystore.password": "***",
            "ssl.keystore.location": "myp12file.p12",
            "auto.offset.reset": "earliest"
        },
        # 其他参数保持不变
      )
    

方案二:在工作节点启动时主动下载证书

如果需要自定义本地存储路径或更灵活的控制,可以通过自定义DoFn,在每个工作节点初始化阶段从GCS下载证书。

  1. 编写下载证书的DoFn
    这个DoFn会在工作节点启动时执行下载操作,确保证书存在于工作节点本地:

    import apache_beam as beam
    from google.cloud import storage
    
    class DownloadSSLCerts(beam.DoFn):
        def __init__(self, gcs_bucket, gcs_file_path, local_path):
            self.gcs_bucket = gcs_bucket
            self.gcs_file_path = gcs_file_path
            self.local_path = local_path
    
        def setup(self):
            # 工作节点初始化时执行下载
            storage_client = storage.Client()
            bucket = storage_client.get_bucket(self.gcs_bucket)
            blob = bucket.blob(self.gcs_file_path)
            blob.download_to_filename(self.local_path)
    
  2. 在Pipeline开头添加下载步骤
    确保证书在Kafka读取操作前完成下载:

    with beam.Pipeline(options=pipeline_options) as pipeline:
        # 先下载SSL证书到工作节点的/tmp目录
        pipeline | "Download SSL Certs" >> beam.Create([None]) | beam.ParDo(DownloadSSLCerts(
            gcs_bucket="你的桶名",
            gcs_file_path="myp12file.p12",
            local_path="/tmp/myp12file.p12"
        ))
    
        # 执行Kafka读取操作
        msg_kv_bytes = pipeline | "read file" >> ReadFromKafka(
             consumer_config={
              # 其他配置不变
              "ssl.truststore.location": "/tmp/myp12file.p12",
              "ssl.keystore.location": "/tmp/myp12file.p12",
              # 其他配置不变
            },
            # 其他参数保持不变
          )
    

关键注意事项

  • 确保Dataflow使用的服务账号拥有目标GCS桶的Storage Object Viewer权限,否则无法读取或下载证书文件。
  • 如果使用本地运行的expansion service,需要确保该服务所在环境也能访问到证书文件;若使用GCP托管的expansion service,Staging方案会自动处理同步。
  • 不要在本地代码中直接执行文件复制到/tmp的操作,这类操作只会在提交任务的本地服务器执行,不会影响Dataflow工作节点。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 08:51:00