【问题标题】:Is there a Kinesis Connector for PyFlink?PyFlink 有 Kinesis 连接器吗?
【发布时间】:2021-03-08 13:56:37
【问题描述】:

我开始研究流式应用程序,并试图弄清楚 PyFlink 是否符合我的要求。我需要能够从 Kinesis Stream 中读取数据。我在the docs 看到有一个 Kinesis Stream 连接器,但我不知道它是否也适用于 Python 版本,如果是,如何配置它。

更新:

我找到了this other doc page,它解释了如何使用 Python 中默认连接器以外的连接器。然后我从here 下载了 Kinesis jar。我下载的版本是flink-connector-kinesis_2.11-1.11.2,与引用的here匹配。

然后,我从文档中的脚本中更改了这一行:t_env.get_config().get_configuration().set_string("pipeline.jars", "file://<absolute_path_to_jar>/connector.jar")

但是,在尝试执行脚本时,我收到了这个 Java 错误:Caused by: org.apache.flink.table.api.ValidationException: Could not find any factory for identifier 'kinesis' that implements 'org.apache.flink.table.factories.DynamicTableSourceFactory' in the classpath.

我也尝试从脚本中删除该配置行,然后以 ./bin/flink run -py <my_script>.py -j ./<path_to_jar>/connector.jar 运行它,但这让我遇到了同样的错误。

我的理解是我添加的 Jar 没有被 Flink 正确识别。我在这里做错了吗?

【问题讨论】:

  • 这能回答你的问题吗? Consuming a kinesis stream in python
  • 我不这么认为,这解释了如何使用 boto 或 Kinesis 客户端库在 Python 中连接到 Kinesis,但我对 PyFlink 的 Kinesis 连接器感兴趣,就像链接中的那个我输入了问题的描述,除了那个只解释了如何在 Scala 和 Java 中使用它,而不是 Python。
  • 不,这不是解释如何仅与 boto 连接。如果您在那里查看另一个答案,您会看到他们还建议使用“Kinesis Client Library (KCL) for Python”
  • 是的,我编辑了我的评论。 KCL 没有回答我的问题,我正在寻找 PyFlink 连接器

标签: python apache-flink


【解决方案1】:

因为要支持的连接器很多,我们需要一个个回馈社区。我们在本地开发了 Kinesis 连接器。由于用户对 Kinesis 连接器有需求,我们将其贡献给 PyFlink。现在PyFlink数据流的相关documentation还在完善中,大家可以先看看Jira看看支持的功能

【讨论】:

    【解决方案2】:

    可能需要澄清一下 PyFlink 当前(Flink 1.11)是 Flink 的 Table API/SQL 的包装器。您尝试使用的连接器是 DataStream API 连接器。

    在 Flink 1.12 中,在接下来的几周内,也会有一个用于 Table API/SQL 的 Kinesis 连接器,所以你应该可以使用它。有关当前支持的连接器的概述,this 是您应该参考的文档页面。

    注意: 正如星博所说,PyFlink 将从 Flink 1.12 开始封装 DataStream API,因此如果您需要较低级别的抽象来实现更复杂的实现,您也可以从 Kinesis 中使用在那里。

    【讨论】:

      猜你喜欢
      • 2022-07-19
      • 1970-01-01
      • 2021-12-24
      • 1970-01-01
      • 1970-01-01
      • 2016-08-20
      • 1970-01-01
      • 2020-11-07
      • 2022-12-19
      相关资源
      最近更新 更多