Apache Beam Python SDK与Firestore依赖冲突及运行报错咨询
问题与解决方案
环境与依赖
- Python版本:3.10.11
requirements.txt内容:
apache-beam==2.46.0 cachetools==4.2.4 certifi==2022.12.7 charset-normalizer==3.1.0 cloudpickle==2.2.1 crcmod==1.7 dill==0.3.1.1 docopt==0.6.2 fastavro==1.7.3 fasteners==0.18 google-api-core==1.34.0 google-apitools==0.5.31 google-auth==2.17.3 google-auth-httplib2==0.1.0 google-cloud-bigquery==3.3.3 google-cloud-bigquery-storage==2.16.0 google-cloud-bigtable==1.7.3 google-cloud-core==1.7.3 google-cloud-datastore==1.15.5 google-cloud-dlp==3.9.0 google-cloud-firestore==2.0.0 google-cloud-language==1.3.2 google-cloud-pubsub==2.13.7 google-cloud-pubsublite==1.8.1 google-cloud-recommendations-ai==0.7.1 google-cloud-spanner==3.22.0 google-cloud-videointelligence==1.16.3 google-cloud-vision==3.1.2 google-crc32c==1.5.0 google-resumable-media==2.5.0 googleapis-common-protos==1.56.4 grpc-google-iam-v1==0.12.4 grpcio==1.54.0 grpcio-status==1.48.2 hdfs==2.7.0 httplib2==0.21.0 idna==3.4 numpy==1.24.3 oauth2client==4.1.3 objsize==0.6.1 orjson==3.8.11 overrides==6.5.0 packaging==21.3 proto-plus==1.22.2 protobuf==3.19.4 pyarrow==9.0.0 pyasn1==0.5.0 pyasn1-modules==0.3.0 pydot==1.4.2 pymongo==3.13.0 pyparsing==3.0.9 python-dateutil==2.8.2 pytz==2023.3 regex==2023.3.23 requests==2.29.0 rsa==4.9 six==1.16.0 sqlparse==0.4.4 typing_extensions==4.5.0 urllib3==1.26.15 zstandard==0.21.0
问题代码
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, GoogleCloudOptions, StandardOptions from google.cloud import firestore from google.api_core.datetime_helpers import DatetimeWithNanoseconds import json import argparse from my_module import get_collections #gets list of FS collections as string class FirestoreToGcs(beam.DoFn): def process(self, element): client = firestore.Client() docs = client.collection(element).get() for doc in docs: doc_dict = doc.to_dict() for key in doc_dict: if isinstance(doc_dict[key], DatetimeWithNanoseconds): doc_dict[key] = doc_dict[key].isoformat() yield json.dumps(doc_dict) def run(argv=None): arg_parser = argparse.ArgumentParser() arg_parser.add_argument( "--output-bucket", dest="output_bucket", required=True, help="GCS bucket to write output files to", ) known_args, pipeline_args = arg_parser.parse_known_args(argv) pipeline_options = PipelineOptions(pipeline_args) with beam.Pipeline(options=pipeline_options) as p: collections = get_collections() for collection in collections: (p | f"Read_{collection}" >> beam.Create([collection]) | f"MapDocuments_{collection}" >> beam.ParDo(FirestoreToGcs()) | f"WriteToGcs_{collection}" >> beam.io.WriteToText(f"{known_args.output_bucket}/{collection}.jsonl") ) if __name__ == '__main__': run()
错误信息
fs_beam_exporter.py", line 13, in process NameError: name 'firestore' is not defined [while running 'MapDocuments_doc-ptransform-46']
用户疑问
发现Stack Overflow上有相关问题指出Apache Beam与其他GCP库存在依赖冲突,请问这些冲突是否已得到解决?
解决方案
1. 解决NameError问题
这个错误和依赖冲突无关,是Beam的序列化机制导致外部导入的模块无法在DoFn内部直接访问,有两种快速修复方法:
- 方法一:在
process方法内部重新导入firestore模块
class FirestoreToGcs(beam.DoFn): def process(self, element): from google.cloud import firestore # 内部导入,避免序列化问题 client = firestore.Client() # 后续代码不变
- 方法二:将firestore客户端作为DoFn实例属性初始化
class FirestoreToGcs(beam.DoFn): def __init__(self): self.client = firestore.Client() def process(self, element): docs = self.client.collection(element).get() # 后续代码不变
2. 关于Apache Beam与GCP库的依赖冲突
你提到的冲突主要集中在protobuf、google-api-core等核心依赖的版本不兼容上,针对你使用的apache-beam==2.46.0:
- 该版本已适配主流GCP客户端库版本,官方修复了大部分常见冲突。从你的requirements.txt来看,
protobuf==3.19.4、google-api-core==1.34.0均符合Beam 2.46.0的依赖要求,不会出现严重冲突。 - 若后续遇到版本兼容问题,建议使用Beam官方提供的
apache-beam[gcp]安装包,它会自动处理所有GCP相关依赖的版本匹配:
pip install apache-beam[gcp]==2.46.0
- 除非有特殊需求,否则不要手动指定GCP客户端库版本,让Beam的依赖管理自动处理版本匹配。
内容的提问来源于stack exchange,提问作者Hossmeister
相关产品推荐
相关产品推荐

