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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 06:27:13