【发布时间】: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