【问题标题】:Flink Python Datastream API Kafka ConsumerFlink Python 数据流 API Kafka 消费者
【发布时间】:2022-04-19 06:41:21
【问题描述】:

我是 pyflink 的新手。我正在尝试编写一个 python 程序来从 kafka 主题中读取数据并将数据打印到标准输出。我点击了链接Flink Python Datastream API Kafka Producer Sink Serializaion。但是由于版本不匹配,我一直看到 NoSuchMethodError。我在https://repo.maven.apache.org/maven2/org/apache/flink/flink-sql-connector-kafka_2.11/1.13.0/flink-sql-connector-kafka_2.11-1.13.0.jar 添加了flink-sql-kafka-connector。有人可以帮我举一个合适的例子吗?以下是我的代码

import json
import os

from pyflink.common import SimpleStringSchema
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.datastream.connectors import FlinkKafkaConsumer
from pyflink.common.typeinfo import Types


def my_map(obj):
    json_obj = json.loads(json.loads(obj))
    return json.dumps(json_obj["name"])


def kafkaread():
    env = StreamExecutionEnvironment.get_execution_environment()

    env.add_jars("file:///automation/flink/flink-sql-connector-kafka_2.11-1.10.1.jar")

    deserialization_schema = SimpleStringSchema()

    kafkaSource = FlinkKafkaConsumer(
        topics='test',
        deserialization_schema=deserialization_schema,
        properties={'bootstrap.servers': '10.234.175.22:9092', 'group.id': 'test'}
    )

    ds = env.add_source(kafkaSource).print()
    env.execute('kafkaread')


if __name__ == '__main__':
    kafkaread()

但是python不识别jar文件并抛出以下错误。

Traceback (most recent call last):
  File "flinkKafka.py", line 31, in <module>
    kafkaread()
  File "flinkKafka.py", line 20, in kafkaread
    kafkaSource = FlinkKafkaConsumer(
  File "/automation/flink/venv/lib/python3.8/site-packages/pyflink/datastream/connectors.py", line 186, in __init__
    j_flink_kafka_consumer = _get_kafka_consumer(topics, properties, deserialization_schema,
  File "/automation/flink/venv/lib/python3.8/site-packages/pyflink/datastream/connectors.py", line 336, in _get_kafka_consumer
    j_flink_kafka_consumer = j_consumer_clz(topics,
  File "/automation/flink/venv/lib/python3.8/site-packages/pyflink/util/exceptions.py", line 185, in wrapped_call
    raise TypeError(
TypeError: Could not found the Java class 'org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer'. The Java dependencies could be specified via command line argument '--jarfile' or the config option 'pipeline.jars'
           

添加jar文件的正确位置是什么?

【问题讨论】:

  • 您使用的是 maven 构建吗?或任何其他构建?

标签: python apache-kafka apache-flink stream-processing data-stream


【解决方案1】:

我看到你下载了flink-sql-connector-kafka_2.11-1.13.0.jar,但是代码加载了flink-sql-connector-kafka_2.11-1.10.1.jar。

也许你可以检查一下

【讨论】:

    【解决方案2】:

    只需要检查 flink-sql-connector jar 的路径

    【讨论】:

      猜你喜欢
      • 2017-10-16
      • 2017-05-31
      • 2017-08-18
      • 2018-10-15
      • 1970-01-01
      • 2022-06-10
      • 1970-01-01
      • 1970-01-01
      • 2018-12-31
      相关资源
      最近更新 更多