能否用AvroKeyValueSinkWriter写入JsonNode?求替代方案
首先明确说一句:不能直接用JsonNode通过AvroKeyValueSinkWriter写入数据,你遇到的ClassCastException就是核心原因——Avro的AvroKeyValueSinkWriter要求值类型必须是IndexedRecord接口的实现类(比如Avro自带的GenericRecord),而JsonNode完全不属于这个类型体系,自然会抛出转换异常。
解决思路:把JSON数据转换成Avro的GenericRecord
你需要把JSON格式的数据(不管是JsonNode还是JSON字符串)转换成Avro规范的GenericRecord对象,再传给AvroKeyValueSinkWriter。下面分步骤给你具体实现方案:
1. 先解析你的Value Schema为Avro的Schema对象
你之前的代码只是把schema字符串存在了properties里,但转换数据时需要用到解析后的Schema实例,所以先补充这一步:
// 在getWriter方法里,解析valueSchema字符串 Schema valueSchemaObj = new Schema.Parser().parse(valueSchema);
2. 把JSON数据转换成GenericRecord
这里有两种常用方式:
方式一:从JsonNode手动转换
如果你已经有了JsonNode对象,可以遍历它的字段,对应到Avro Schema的字段来构建GenericRecord:
private static GenericRecord jsonNodeToGenericRecord(JsonNode jsonNode, Schema schema) { GenericRecord record = new GenericData.Record(schema); for (Schema.Field field : schema.getFields()) { String fieldName = field.name(); JsonNode fieldNode = jsonNode.get(fieldName); if (fieldNode == null || fieldNode.isNull()) { record.put(fieldName, null); continue; } // 根据字段类型做对应转换 switch (field.schema().getType()) { case INT: record.put(fieldName, fieldNode.asInt()); break; case STRING: record.put(fieldName, fieldNode.asText()); break; // 其他类型(比如FLOAT、BOOLEAN等)可以按需补充 default: throw new IllegalArgumentException("Unsupported field type: " + field.schema().getType()); } } return record; }
方式二:直接从JSON字符串解析成GenericRecord(更高效)
如果你的原始数据是JSON字符串,完全可以跳过JsonNode这一步,用Avro自带的JsonDecoder直接解析成GenericRecord,省去中间转换:
private static GenericRecord jsonStringToGenericRecord(String jsonString, Schema schema) throws IOException { GenericRecord record = new GenericData.Record(schema); Decoder decoder = DecoderFactory.get().jsonDecoder(schema, jsonString); new GenericDatumReader<GenericRecord>(schema).read(record, decoder); return record; }
3. 修改测试用例,使用GenericRecord写入
把原来的JsonNode替换成转换后的GenericRecord即可:
@Test public void testWriter() throws Exception { AvroKeyValueSinkWriter<String, GenericRecord> writer = getWriter(); String key = "5639281840180123"; String jsonString = "{\n" + " \"InvoiceNo\": 5370812,\n" + " \"StockCode\": \"22409\",\n" + " \"StoreID\": 0,\n" + " \"TransactionID\": \"537081210180130\"\n" + "}"; // 先解析value schema Schema valueSchemaObj = new Schema.Parser().parse(valueSchema); // 转换为GenericRecord GenericRecord record = jsonStringToGenericRecord(jsonString, valueSchemaObj); writer.setSyncOnFlush(true); writer.open(fs, path); writer.write(new Tuple2<String, GenericRecord>(key, record)); writer.flush(); // 这里可以添加断言验证写入结果 }
额外优化:用Avro生成的实体类替代GenericRecord
如果你的schema固定,还可以用Avro工具(比如avro-tools)根据你的value schema生成对应的Java实体类,这个类会自动实现IndexedRecord接口,使用起来更类型安全,不需要手动处理字段转换。比如生成Transaction类后,直接用ObjectMapper把JSON字符串转成Transaction对象再写入即可。
内容的提问来源于stack exchange,提问作者Chris Snow

