Dataflow Python Flex模板启动失败:提示需安装Java
问题:Dataflow Flex模板启动时提示Java未安装(已在Dockerfile中安装Java)
我正在运行一个从PubSub Lite到BigQuery的Dataflow Flex模板作业,已经在Dockerfile中安装了openjdk-11-jdk并设置了JAVA_HOME,且验证镜像内Java可正常访问,但作业启动时仍报错提示Java未安装。
相关代码与配置
Python 代码
from __future__ import annotations import argparse import json import logging import apache_beam.io.gcp.pubsublite as psub_lite import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions # Defines the BigQuery schema for the output table. schema = 'trip_id:INTEGER,vendor_id:INTEGER,trip_distance:FLOAT,fare_amount:STRING,store_and_fwd_flag:STRING' class ModifyDataForBQ(beam.DoFn): def process(self, pubsub_message, *args, **kwargs): # attributes = dict(pubsub_message.attributes) obj = json.loads(pubsub_message.message.data.decode("utf-8")) yield obj def run( subscription_id: str, dataset: str, table: str, beam_args: list[str] = None, ) -> None: options = PipelineOptions(beam_args, save_main_session=True, streaming=True) table = '{}.{}'.format(dataset, table) p = beam.Pipeline(options=options) pubsub_pipeline = ( p | 'Read from pubsub lite topic' >> psub_lite.ReadFromPubSubLite(subscription_path=subscription_id) | 'Print Message' >> beam.ParDo(ModifyDataForBQ()) | 'Write Record to BigQuery' >> beam.io.WriteToBigQuery(table=table, schema=schema, write_disposition=beam.io.BigQueryDisposition .WRITE_APPEND, create_disposition=beam.io.BigQueryDisposition .CREATE_IF_NEEDED, ) ) result = p.run() result.wait_until_finish() if __name__ == "__main__": logging.getLogger().setLevel(logging.INFO) parser = argparse.ArgumentParser() parser.add_argument( "--subscription_id", type=str, help="Region of Pub/Sub Lite subscription.", default=None ) parser.add_argument( "--dataset", type=str, help="BigQuery Dataset name.", default=None ) parser.add_argument( "--table", type=str, help="BigQuery destination table name.", default=None ) args, beam_args = parser.parse_known_args() run( subscription_id=args.subscription_id, dataset=args.dataset, table=args.table, beam_args=beam_args, )
Dockerfile
FROM gcr.io/dataflow-templates-base/python3-template-launcher-base ENV FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE="/template/requirements.txt" ENV FLEX_TEMPLATE_PYTHON_PY_FILE="/template/streaming_beam.py" COPY . /template RUN apt-get update \ && apt-get install -y openjdk-11-jdk libffi-dev git \ && rm -rf /var/lib/apt/lists/* \ # Upgrade pip and install the requirements. && pip install --no-cache-dir --upgrade pip \ && pip install --no-cache-dir -r $FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE \ # Download the requirements to speed up launching the Dataflow job. && pip download --no-cache-dir --dest /tmp/dataflow-requirements-cache -r $FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE ENV JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64 ENV PIP_NO_DEPS=True ENTRYPOINT ["/opt/google/dataflow/python_template_launcher"]
模板构建命令
gcloud dataflow flex-template build gs://my-bucket-xxxx/templates/streaming-beam-sql.json \ --image-gcr-path "us-central1-docker.pkg.dev/xxxx-xxx-2/dataflow-pubsublite-bigquery/test:latest" \ --sdk-language "PYTHON" \ --flex-template-base-image "PYTHON3" \ --metadata-file "metadata.json" \ --py-path "." \ --env "FLEX_TEMPLATE_PYTHON_PY_FILE=streaming_beam.py" \ --env "FLEX_TEMPLATE_PYTHON_REQUIREMENTS_FILE=requirements.txt" \ --project "xxxx-xxx-2"
模板调用命令
gcloud dataflow flex-template run "streaming-beam-sql" \ --template-file-gcs-location gs://my-bucket-xxxx/templates/streaming-beam-sql.json \ --project "xxxx-xxx-2" \ --parameters "subscription_id=projects/xxxx-xxx-/locations/us-central1/subscriptions/data-streaming-xxxx-subscription,dataset=omer_poc,table=trip2"
错误日志
INFO 2023-06-08T22:27:23.260235Z INFO:root:Starting a JAR-based expansion service from JAR /root/.apache_beam/cache/jars/beam-sdks-java-io-google-cloud-platform-expansion-service-2.41.0.jar INFO 2023-06-08T22:27:23.261209Z ERROR:apache_beam.utils.subprocess_server:Error bringing up service INFO 2023-06-08T22:27:23.261252Z Traceback (most recent call last): INFO 2023-06-08T22:27:23.261270Z File "/usr/local/lib/python3.7/site-packages/apache_beam/utils/subprocess_server.py", line 79, in start INFO 2023-06-08T22:27:23.261296Z endpoint = self.start_process() INFO 2023-06-08T22:27:23.261313Z File "/usr/local/lib/python3.7/site-packages/apache_beam/utils/subprocess_server.py", line 181, in start_process INFO 2023-06-08T22:27:23.261329Z 'Java must be installed on this system to use this ' INFO 2023-06-08T22:27:23.261343Z RuntimeError: Java must be installed on this system to use this transform/runner.
解决方案
1. 确保Java可执行文件在系统PATH中
Apache Beam的subprocess_server模块会直接调用java命令,若PATH中未包含Java的bin目录,会导致找不到Java。修改Dockerfile,在设置JAVA_HOME后添加PATH配置:
ENV JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64 ENV PATH="$JAVA_HOME/bin:$PATH" # 新增该行
2. 验证镜像中的Java环境
重新构建镜像后,运行以下命令确认Java可正常调用:
docker run --rm <your-image-id> java -version
若能输出Java版本信息,说明环境配置正确。
3. 检查Apache Beam版本兼容性
错误中使用的beam版本为2.41.0,部分旧版本可能存在Java路径检测的问题。尝试升级requirements.txt中的apache-beam版本到较新的稳定版(如2.50.0+),重新构建模板镜像。
内容的提问来源于stack exchange,提问作者danny.lesnik
相关产品推荐
相关产品推荐

