【问题标题】:Using Python Processors in Java Flink Application在 Java Flink 应用程序中使用 Python 处理器
【发布时间】:2021-06-17 20:52:52
【问题描述】:

我有一个用例,我想用 Java 中的 Flink 实现 AWS Kinesis Data 应用程序。它将通过 Data Streams API 监听多个 Kinesis 流。但是,这些流的分析将在 Python 中完成(因为我们的数据科学家更喜欢 Python)。

来自this answer,似乎支持从 Java 调用 Python UDF。但是,我希望能够将传入流转换为表格,通过

StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
Table sessionsTable = tableEnv.fromDataStream(inputStream);

...然后有一个 Python 处理器被调用来处理该流。

我真的有3个问题:

  1. 这是受支持的用例吗?
  2. 如果有,是否有说明如何操作的文档?
  3. 如果是这样,这是否会显着增加应用程序的开销?

【问题讨论】:

    标签: apache-flink pyflink amazon-kinesis-analytics


    【解决方案1】:

    Flink 文档中学习如何将 Python 与表和数据流结合使用的起点位于 https://ci.apache.org/projects/flink/flink-docs-stable/docs/dev/python/overview/

    Python API 仅提供 Java 的一部分;您必须查看是否包含您需要的内容。

    不确定性能,但您可以,例如,convert back and forth between Flink Tables and Pandas dataframes

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2012-09-01
      • 2021-03-15
      • 2021-08-09
      • 1970-01-01
      相关资源
      最近更新 更多