Apache Beam Dataflow运行时Worker无法找到MongoDB证书文件
我编写了一个基于Apache Beam从MongoDB读取数据的函数,代码如下:
def create_mongo_pipeline(p: beam.Pipeline, mongo_uri: str, db: str, coll: str, cert_file: str, gcs_bucket: str) -> None: # Read from MongoDB docs = p | 'ReadMyFile' >> beam.io.mongodbio.ReadFromMongoDB( uri=mongo_uri, db=db, coll=coll, bucket_auto=True, extra_client_params={ "tls": True, "authMechanism": "MONGODB-X509", "tlsCertificateKeyFile": cert_file, } ) ...
证书文件路径定义在config.json中:
{ "connection": { "type": "mongo", "uri": "mongodb://blablabla:27017", "db": "mydb", "certificate": "resources/mycert.pem" } }
本地使用Direct Runner运行完全正常,但切换到Dataflow Runner时,出现如下错误:
RuntimeError: FileNotFoundError: [Errno 2] No such file or directory: 'resources/mycert.pem' [while running 'ReadMyFile/Read/SDFBoundedSourceReader/ParDo(SDFBoundedSourceDoFn)/PairWithRestriction-ptransform-37']
我尝试将代码打包进Docker容器,仍然遇到相同的错误。调用Dataflow Runner的代码如下:
beam_options = PipelineOptions( runner="DataflowRunner", project=config["project_id"], job_name=config["job_name"], experiments=["use_grpc_for_gcs"], temp_location=f"{config['gcs_bucket']}/temp", )
Dataflow Runner运行时,Worker节点无法直接访问本地文件系统或Docker容器内的文件路径,核心原因是:
- Direct Runner在本地执行,能读取本地的
resources/mycert.pem;但Dataflow的Worker是GCP云端虚拟机,本地路径对它们完全不可见。 - 即使打包进Docker,Worker启动时不会自动同步容器内的证书文件,且MongoDB客户端会尝试从Worker的本地文件系统读取路径,而非容器内路径。
正确处理步骤
将证书文件上传到GCS存储桶
把resources/mycert.pem上传到你的GCS桶中,比如路径gs://your-bucket/certs/mycert.pem。修改代码,在Worker运行时下载证书到本地
在读取MongoDB之前,添加步骤将GCS上的证书下载到Worker的临时目录,再用这个临时路径作为证书文件路径:import tempfile from google.cloud import storage def download_cert_from_gcs(gcs_cert_path: str) -> str: # 解析GCS路径,分离桶名和对象路径 bucket_name, blob_path = gcs_cert_path.replace("gs://", "").split("/", 1) storage_client = storage.Client() bucket = storage_client.bucket(bucket_name) blob = bucket.blob(blob_path) # 创建临时文件保存证书 with tempfile.NamedTemporaryFile(mode='wb', delete=False, suffix='.pem') as temp_file: blob.download_to_file(temp_file) return temp_file.name def create_mongo_pipeline(p: beam.Pipeline, mongo_uri: str, db: str, coll: str, gcs_cert_path: str, gcs_bucket: str) -> None: # 先下载证书到Worker本地 cert_file = download_cert_from_gcs(gcs_cert_path) # Read from MongoDB docs = p | 'ReadMyFile' >> beam.io.mongodbio.ReadFromMongoDB( uri=mongo_uri, db=db, coll=coll, bucket_auto=True, extra_client_params={ "tls": True, "authMechanism": "MONGODB-X509", "tlsCertificateKeyFile": cert_file, } ) ...更新config.json中的证书路径
把certificate的值改成GCS路径:{ "connection": { "type": "mongo", "uri": "mongodb://blablabla:27017", "db": "mydb", "certificate": "gs://your-bucket/certs/mycert.pem" } }确保Dataflow Worker有GCS访问权限
运行Dataflow的服务账号需要拥有GCS存储桶的storage.objects.get权限,否则无法下载证书文件。
备选方案:使用Dataflow文件Staging
也可以通过files_to_stage参数把证书文件打包到Worker环境中,示例:
beam_options = PipelineOptions( runner="DataflowRunner", project=config["project_id"], job_name=config["job_name"], experiments=["use_grpc_for_gcs"], temp_location=f"{config['gcs_bucket']}/temp", files_to_stage=["resources/mycert.pem"] )
此时证书会被放到Worker的工作目录,你需要获取该目录的绝对路径来设置tlsCertificateKeyFile,但这种方式灵活性不如GCS下载,推荐优先使用前者。
内容的提问来源于stack exchange,提问作者RDGuida

