【问题标题】:Error while using KafkaIO in apache beam DirectRunner在 Apache Beam DirectRunner 中使用 KafkaIO 时出错
【发布时间】:2020-07-07 19:22:07
【问题描述】:

我正在使用 apache beam DirectRunner 从 kafka 主题加载数据。我的代码如下:

conf={'bootstrap.servers':'localhost:9092'}

with beam.Pipeline() as pipeline:
        (pipeline
        |       ReadFromKafka(consumer_config=conf,topics=['topic1'])
        )

我正在使用以下命令来运行此代码:

python3 topic_to_gcs --runner DirectRunner

出现以下错误:

File "/usr/lib/python3.7/subprocess.py", line 1522, in _execute_child
    raise child_exception_type(errno_num, err_msg, err_filename)
FileNotFoundError: [Errno 2] No such file or directory: 'docker': 'docker'

提前致谢:)

【问题讨论】:

  • 你是在 docker 容器中运行这个吗?
  • @bigbounty,否...在 gcp 计算实例(基础机器)上。

标签: python ubuntu apache-kafka apache-beam apache-beam-io


【解决方案1】:

目前,Apache Beam 在 Python SDK 中使用所谓的外部转换从 Kafka 读取数据。这实际上意味着,您的 Python 管道将生成一个 Java 容器并从容器内部连接到 Kafka。然后它将数据传回你的 Python 管道(更多关于这个here)。

如果您可以在运行管道的主机上安装 docker(以及在您计划运行它的所有其他位置,如果您将运行器从 DirectRunner 更改为某个分布式运行器),那么这将是最佳选择去。

否则您可以在我的回答here 中了解当前状态。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-03-22
    • 1970-01-01
    • 2020-03-07
    • 2022-12-24
    • 1970-01-01
    相关资源
    最近更新 更多