【发布时间】:2019-11-05 19:23:07
【问题描述】:
我希望获得有关在运行 Beam wordcount.py 演示时如何设置 --environment_config 的指导。
它与 DirectRunner 一起运行良好。 Flink 的 wordcount 也运行良好(即通过flink run 运行 Flink)。
我想使用 Flink 运行器运行 Beam,使用 beam documentation 中描述的“单独的 Flink 集群”。我用不了Docker,所以打算用--environment_type=PROCESS。
我在 python 代码中使用以下内容来设置 environment_config:
environment_config = dict()
environment_config['os'] = platform.system().lower()
environment_config['arch'] = platform.machine()
environment_config['command'] = 'ls'
ec = "--environment_config={}".format(json.dumps(environment_config))
显然命令不正确。当我运行它时,Flink 确实接收并成功处理了 DataSource 子任务。它最终在CHAIN MapPartitions 上超时。
有人可以提供有关如何设置 environment_config 的指导(或链接)吗?我在 Singularity 容器中运行 Beam。
【问题讨论】: