【问题标题】:flink with python, execution of job failedflink与python,作业执行失败
【发布时间】:2019-03-06 13:36:14
【问题描述】:

第一次尝试时,我想从文件中读取 JSON 数据并将其传递给 Flink。我定义了一个源(逐行读取 JSON 字符串)和一个占位符过滤器。见代码:

from org.apache.flink.streaming.api.functions.source import SourceFunction
from org.apache.flink.api.common.functions import FilterFunction
import json
import sys

class Json_reader(SourceFunction):
    def readjason(self, ctx):
        sys.stdin = open('capture.json', 'r')
        for line in sys.stdin:
            ctx.collect(json.loads(line))


class Dummy_Filter(FilterFunction):
    def filter(self, value):
        return True

#
# The pipeline definition.
#
def main(factory):
    env = factory.get_execution_environment()
    env.create_python_source(Json_reader()) \
        .filter(Dummy_Filter()) \
        .output()
    env.execute()

当我构建作业并将其移动到我启动的 Flink 集群时,我收到以下错误消息:

VirtualBox:/media/sf_Python$ ./flink-1.7.2/bin/pyflink-stream.sh ./json_parser_flink.py 开始执行程序 运行失败 计划:空回溯(最近一次通话):文件“”,行 1、在文件中 "/tmp/flink_streaming_plan_fbe13c4c-6918-46d4-a4bc-36908a2bea24/json_parser_flink.py", 第 25 行,主要在 org.apache.flink.client.program.rest.RestClusterClient.submitJob(RestClusterClient.java:268) 在 org.apache.flink.client.program.ClusterClient.run(ClusterClient.java:487) 在 org.apache.flink.streaming.api.environment.StreamContextEnvironment.execute(StreamContextEnvironment.java:66) 在 org.apache.flink.streaming.api.environment.StreamExecutionEnvironment.execute(StreamExecutionEnvironment.java:1510) 在 org.apache.flink.streaming.python.api.environment.PythonStreamExecutionEnvironment.execute(PythonStreamExecutionEnvironment.java:245) 在 sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 在 sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) 在 sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 在 java.lang.reflect.Method.invoke(Method.java:498) org.apache.flink.client.program.ProgramInvocationException: org.apache.flink.client.program.ProgramInvocationException:作业 失败的。 (职位编号:31615948194c951be03d46576929aa23)

该程序不包含 Flink 作业。也许你忘了打电话 在执行环境上执行()。

我没有忘记调用 execute()。

【问题讨论】:

    标签: python apache-flink pyflink


    【解决方案1】:

    我发现了问题。 Fast 期望 SourceFunction 中有一个 run() 函数。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-11-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多