如何将JSONObject转换为Avro GenericRecord?Apache Beam实践求助
解决Apache Beam中修改GenericRecord后转换回原类型的问题
方案一:直接操作GenericRecord(推荐)
不需要转成JSONObject,直接基于Avro的API修改GenericRecord和嵌套结构,既高效又避免类型转换问题:
import org.apache.avro.generic.GenericArray import org.apache.avro.generic.GenericRecord import org.apache.avro.generic.GenericData val attributes = dataNestedObjGroup.getAsGenericRecord("details")?.get("attributes") as GenericArray<GenericRecord> attributes.forEach { attribute -> // 遍历当前attribute的所有字段 attribute.schema.fields.forEach { field -> val fieldName = field.name // 安全转换为嵌套的GenericRecord(避免类型异常) val innerRecord = attribute.get(fieldName) as? GenericRecord ?: return@forEach // 检查是否为目标字段 if (innerRecord.get("name") == "customer_id") { // 获取"value"字段的Schema,创建空的GenericArray val valueFieldSchema = innerRecord.schema.getField("value").schema() val emptyValueArray = GenericData.Array(0, valueFieldSchema) // 直接更新嵌套记录的value字段 innerRecord.put("value", emptyValueArray) } } }
方案二:将JSONObject转回GenericRecord
如果必须通过JSON中转,可以利用Avro的JSON解码器完成转换,核心是依赖原GenericRecord的Schema:
import org.apache.avro.Schema import org.apache.avro.generic.GenericDatumReader import org.apache.avro.io.DecoderFactory import org.apache.avro.generic.GenericRecord import org.apache.avro.generic.GenericArray val attributes = dataNestedObjGroup.getAsGenericRecord("details")?.get("attributes") as GenericArray<GenericRecord> // 用索引遍历,方便替换数组元素 for (i in attributes.indices) { val attribute = attributes[i] val jsonRootObject = JSONObject(attribute.toString()) jsonRootObject.keySet().forEach { keyStr -> val value = jsonRootObject.get(keyStr) val innerObject = JSONObject(value?.toString()) val valueOfKey = innerObject.getJSONArray("value") if (innerObject.get("name") == "customer_id" && valueOfKey != null) { innerObject.put("value", JsonArray()) } jsonRootObject.put(keyStr, innerObject) } // 利用原attribute的Schema将JSON转回GenericRecord val attributeSchema = attribute.schema val jsonString = jsonRootObject.toString() val reader = GenericDatumReader<GenericRecord>(attributeSchema) val decoder = DecoderFactory.get().jsonDecoder(attributeSchema, jsonString) val updatedRecord = reader.read(null, decoder) // 替换数组中的原记录 attributes[i] = updatedRecord }
注意事项
- GenericRecord本身是可变类型,修改嵌套的子记录会直接同步到原对象中,无需额外替换
- 操作时建议使用
as?进行安全类型转换,避免ClassCastException - 若GenericArray是不可变实现(极少数情况),可以创建新的
GenericData.Array,将修改后的记录逐个加入后替换原数组
内容的提问来源于stack exchange,提问作者Kulasangar
相关产品推荐
相关产品推荐

