【发布时间】:2022-01-06 12:52:30
【问题描述】:
我是 Flink 流媒体的初学者。
使用 RowCsvInputFormat 读取文件时,Kryo 序列化程序创建 Row 的代码无法正常工作。
代码如下。
val readLocalCsvFile = new RowCsvInputFormat(
new Path("flink-test/000000_1"),
Array(Types.STRING, Types.STRING, Types.STRING),
"\n",
","
)
val read = env.readFile(
readLocalCsvFile,
"flink-test/000000_1",
FileProcessingMode.PROCESS_CONTINUOUSLY,
1000000)
read.print()
env.execute("test")
文件000000_1的内容如下。
aa,bb,cc
aaa,bbb,ccc
经过调试,我很好地得到了aa、bb、cc的划分值。但是当我将这些值一一放入Row的字段时,由于字段为空,引发了一个nullpointexception。
下图显示 Row 的字段为空。
上述代码执行时创建Row的代码如下。 KryoSerializer 生成行。
val kryo = new EmptyFlinkScalaKryoInstantiator().newKryo
val Row = kryo.newInstance(classOf[Row])
输出错误如下。
java.lang.NullPointerException
at org.apache.flink.types.Row.setField(Row.java:140)
at org.apache.flink.api.java.io.RowCsvInputFormat.fillRecord(RowCsvInputFormat.java:162)
at org.apache.flink.api.java.io.RowCsvInputFormat.fillRecord(RowCsvInputFormat.java:33)
at org.apache.flink.api.java.io.CsvInputFormat.readRecord(CsvInputFormat.java:113)
at org.apache.flink.api.common.io.DelimitedInputFormat.nextRecord(DelimitedInputFormat.java:551)
at org.apache.flink.api.java.io.CsvInputFormat.nextRecord(CsvInputFormat.java:80)
at org.apache.flink.streaming.api.functions.source.ContinuousFileReaderOperator.readAndCollectRecord(ContinuousFileReaderOperator.java:387)
at
【问题讨论】:
标签: scala file streaming apache-flink kryo