You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 17:10:02