【发布时间】:2019-04-30 05:20:59
【问题描述】:
我正在尝试使用 spark 结构化流 (version-2.3.1) 处理来自 kafka 的流式 avro 数据,因此我尝试使用 this 示例进行反序列化。
它仅在主题 value 部分包含 StringType 时才有效,但在我的情况下,架构包含 long and integers,如下所示:
public static final String USER_SCHEMA = "{"
+ "\"type\":\"record\","
+ "\"name\":\"variables\","
+ "\"fields\":["
+ " { \"name\":\"time\", \"type\":\"long\" },"
+ " { \"name\":\"thnigId\", \"type\":\"string\" },"
+ " { \"name\":\"controller\", \"type\":\"int\" },"
+ " { \"name\":\"module\", \"type\":\"int\" }"
+ "]}";
所以它给出了一个例外
sparkSession.udf().register("deserialize", (byte[] data) -> {
GenericRecord record = recordInjection.invert(data).get(); //throws error at invert method.
return RowFactory.create(record.get("time"), record.get("thingId").toString(), record.get("controller"), record.get("module"));
}, DataTypes.createStructType(type.fields()));
说
Failed to invert: [B@22a45e7
Caused by java.io.IOException: Invalid int encoding.
因为我在架构 long and int 类型中有 time, controller and module。
我猜这是字节数组byte[] data的某种编码和解码格式错误。
【问题讨论】:
标签: apache-spark apache-spark-sql byte spark-structured-streaming