【问题标题】:Flink throwing com.esotericsoftware.kryo.KryoException: java.lang.NullPointerExceptionFlink 抛出 com.esotericsoftware.kryo.KryoException: java.lang.NullPointerException
【发布时间】:2020-09-11 02:24:45
【问题描述】:

我在我的 flink 流作业中看到一个奇怪的行为。这是我的代码

        streamExecutionEnvironment.enableCheckpointing(checkPointInterval, CheckpointingMode.EXACTLY_ONCE);
        streamExecutionEnvironment.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
        ExecutionConfig executionConfig = streamExecutionEnvironment.getConfig();
        executionConfig.disableForceKryo();
        executionConfig.enableForceAvro();
        Path path = new Path(outputPath);
        CheckpointConfig config = streamExecutionEnvironment.getCheckpointConfig();
        config.enableExternalizedCheckpoints(ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

        String mutateConfig = IOUtils.toString(EventProcessor.class.getClassLoader().getResourceAsStream(configFile));

        FlinkKafkaConsumer flinkKafkaConsumer = new FlinkKafkaConsumer(topics,
                new KafkaGenericAvroDeserializationSchema(schemaRegistryUrl),
                properties);

flinkKafkaConsumer.setCommitOffsetsOnCheckpoints(true);
        DataStream<GenericRecord> dataStream = streamExecutionEnvironment.addSource(flinkKafkaConsumer).name("booking_flow_source");


        DataStream<GenericRecord> enrichDataStream = dataStream.map(new MapFunction<GenericRecord, GenericRecord>() {
            private transient Mutator mutator;
            @Override
            public GenericRecord map(GenericRecord record)  {
                GenericRecord mutateRecord=record;
                try {
                    mutator = new Mutator(mutateConfig);
                    mutateRecord = mutator.mutate(record);
                } catch (Exception e) {
                    e.printStackTrace();
                }
                return mutateRecord;
            }
        });

        enrichDataStream.print();

到目前为止,此代码运行良好。现在我需要从我的 avro 模式生成 java 类,所以我已经包含了这个 avro 依赖项。

<dependency>
                <groupId>org.apache.avro</groupId>
                <artifactId>avro</artifactId>
                <version>1.9.1</version>
</dependency>

在我的 pom 中包含这个之后,我的代码停止工作并且我得到了异常:

org.apache.flink.streaming.runtime.tasks.ExceptionInChainedOperatorException: Could not forward element to next operator

com.esotericsoftware.kryo.KryoException: java.lang.NullPointerException
Serialization trace:
props (org.apache.avro.Schema$Field)
fieldMap (org.apache.avro.Schema$RecordSchema)
schema (org.apache.avro.generic.GenericData$Record)

即使我在我的代码中禁用 kryo 并强制 avro,我仍然得到相同的异常。 如果我删除此依赖项,则代码正在运行并且我的流正在打印。

所以我无法通过添加 avro 依赖项来理解发生了什么变化。

请帮忙

【问题讨论】:

    标签: java apache-kafka apache-flink avro flink-streaming


    【解决方案1】:

    我遇到了类似的问题。我修复了它在 flinkconfiguration 中设置 classloader.resolve-order: parent-first。

    【讨论】:

      猜你喜欢
      • 2021-04-26
      • 1970-01-01
      • 2013-12-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-06-04
      • 1970-01-01
      相关资源
      最近更新 更多