【问题标题】:Apache Flink Python Table API UDF Dependencies ProblemApache Flink Python Table API UDF 依赖问题
【发布时间】:2020-05-05 18:34:31
【问题描述】:

通过将涉及用户定义函数 (UDF) 的 Python 表 API 作业提交到本地集群启动后,它会崩溃并出现 py4j.protocol.Py4JJavaError由

引起

java.util.ServiceConfigurationError:org.apache.beam.sdk.options.PipelineOptionsRegistrar:org.apache.beam.sdk.options.DefaultPipelineOptionsRegistrar 不是子类型。

我知道这是一个关于 lib 路径/类加载依赖项的错误。我已经尝试按照以下链接中的所有说明进行操作:https://ci.apache.org/projects/flink/flink-docs-release-1.10/monitoring/debugging_classloading.html

我已经使用classloader.parent-first-patterns-additional 配置选项尝试了多种不同的配置。带有org.apache.beam.sdk.[...] 的不同条目会导致不同的附加错误消息。

以下依赖于 apache beam 的依赖项位于 lib 路径上:

  • beam-model-fn-execution-2.20.jar
  • beam-model-job-management-2.20.jar
  • beam-model-pipeline-2.20.jar
  • beam-runners-core-construction-java-2.20.jar
  • beam-runners-java-fn-execution-2.20.jar
  • beam-sdks-java-core-2.20.jar
  • beam-sdks-java-fn-execution-2.20.jar
  • beam-vendor-grpc-1_21_0-0.1.jar
  • beam-vendor-grpc-1_26_0.0.3.jar
  • beam-vendor-guava-26_0-jre-0.1.jar
  • beam-vendor-sdks-java-extensions-protobuf-2.20.jar

我也可以排除是我的代码问题,因为我测试了项目网站的如下示例代码:https://flink.apache.org/2020/04/09/pyflink-udf-support-flink.html

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, DataTypes
from pyflink.table.descriptors import Schema, OldCsv, FileSystem
from pyflink.table.udf import udf

env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)
t_env = StreamTableEnvironment.create(env)

add = udf(lambda i, j: i + j, [DataTypes.BIGINT(), DataTypes.BIGINT()], DataTypes.BIGINT())

t_env.register_function("add", add)

t_env.connect(FileSystem().path('/tmp/input')) \
    .with_format(OldCsv()
                 .field('a', DataTypes.BIGINT())
                 .field('b', DataTypes.BIGINT())) \
    .with_schema(Schema()
                 .field('a', DataTypes.BIGINT())
                 .field('b', DataTypes.BIGINT())) \
    .create_temporary_table('mySource')

t_env.connect(FileSystem().path('/tmp/output')) \
    .with_format(OldCsv()
                 .field('sum', DataTypes.BIGINT())) \
    .with_schema(Schema()
                 .field('sum', DataTypes.BIGINT())) \
    .create_temporary_table('mySink')

t_env.from_path('mySource')\
    .select("add(a, b)") \
    .insert_into('mySink')

t_env.execute("tutorial_job")

执行此代码时,会出现相同的错误消息。

谁有 Flink 集群配置的描述,可以使用 UDF 运行 Python Table API 作业?非常感谢您提前提供的所有提示!

【问题讨论】:

    标签: java python apache-flink apache-beam py4j


    【解决方案1】:

    Apache Flink 的新版本 1.10.1 解决了这个问题。现在可以使用命令run -py path/to/script 通过二进制文件执行问题中显示的示例脚本,没有任何问题。

    至于依赖关系,它们已经包含在已经交付的flink_table_x.xx-1.10.1.jar 中。因此,无需将进一步的依赖项添加到 lib-path,这是通过调试/配置尝试在问题中完成的。

    【讨论】:

      猜你喜欢
      • 2016-08-22
      • 1970-01-01
      • 1970-01-01
      • 2019-07-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-02-21
      相关资源
      最近更新 更多