寻求Golang/Python代码:从JSON文档推断AVRO Schema
从JSON文档推断AVRO Schema的实现方案
我尝试编写JSON解析器以生成AVRO Schema,但未成功实现。现寻求可从JSON文档推断AVRO Schema的Golang代码,若有可用的Python代码也可接受。我已编写了一段Python代码雏形(可生成通用对象表示,目标是转为AVRO Schema),并附上了期望得到的AVRO Schema示例:
我的Python代码雏形
import json import urllib TYPES = { type(1): 'long', type(1.2): 'double', type("abc"): 'string', type(u"abc"): 'string', type(True): 'boolean', type([]): 'array', type(()): 'array', type({}): 'object', type(None): 'null', } COMPOUND_TYPES = frozenset(['array', 'object']) SCALARS_TYPES = set(TYPES.values()) - COMPOUND_TYPES def parse_sample(item, paths=None, base=None, avro=None): avro = avro or {} base = base or () paths = paths or {} # path container and counters paths.setdefault(base, 0) paths[base] += 1 type_ = TYPES.get(type(item), "any") base1 = base + (type_,) paths.setdefault(base1, 0) print("ppp--- ", paths) print("bbb+++", base1) print(" tttt ",type_) if type_ in SCALARS_TYPES: paths[base1] += 1 # avro.append({"name": k, "type": TYPES.get(type(item), "any")}) elif type_ == "array": paths[base1] += 1 base1b = base1 + (None,) # adding extra place for possible extensions (eg array index) paths.setdefault(base1b, 0) for subitem in item: parse_sample(subitem, paths=paths, base=base1b) elif type_ == "object": avro.update({ " type": "record", "name": "Record"}) paths[base1] += 1 for (k, subitem) in item.items(): base1b = base1 + (k,) print( subitem, TYPES.get(type(subitem), "any"), paths, base1b) parse_sample(subitem, paths=paths, base=base1b) # print(paths) # print(avro) return paths def from_json(url): u = urllib.urlopen(url).read() return json.loads(u) def guess_schema(s): paths = parse_sample(s) # print(paths) # return build_schema(paths) json_content = { "timestamp": 1661193367, "plant": "M14", "group": "g3", "device": "Sensor3", "parameter": "Power", "value": 100.3000000000001, "unit": "V", "limits": { "class": "A", "high": 200.2, "low": 32 }, "portfolio": { "id": "M00" } } def write(schema, filename): schema_str = json.dumps(schema) # print(schema) # print(json.dumps(schema)) with open(filename, 'w') as x: json.dump(schema, x, indent=4 ) # write(avro_guess(json_content), "guess_avro.json") paths = parse_sample(json_content) # print(paths)
期望生成的AVRO Schema
{ "type": "object", "name": "Record", "fields": { "device": { "required": true, "type": "string" }, "parameter": { "required": true, "type": "string" }, "plant": { "required": true, "type": "string" }, "timestamp": { "required": true, "type": "long" }, "unit": { "required": true, "type": "string" }, "value": { "required": true, "type": "double" }, "group": { "required": true, "type": "string" }, "limits": { "type": "object", "name": "Record", "fields": { "class": { "required": false, "type": "string" }, "high": { "required": false, "type": "double" }, "low": { "required": false, "type": "double" } } }, "portfolio": { "type": "object", "name": "Record", "fields": { "id": { "required": true, "type": "string" } } } } }
完善后的Python实现代码
基于你的雏形修改,以下代码可正确生成符合预期的AVRO Schema:
import json # JSON类型到AVRO类型的映射 TYPE_MAPPING = { int: "long", float: "double", str: "string", bool: "boolean", type(None): "null", list: "array", dict: "record" } def infer_schema(obj, record_name="Record", required=True): """递归推断JSON结构对应的AVRO Schema""" schema = {} obj_type = type(obj) # 处理基础标量类型 if obj_type in TYPE_MAPPING: avro_type = TYPE_MAPPING[obj_type] schema["type"] = avro_type schema["required"] = required return schema # 处理数组类型 if isinstance(obj, list): if not obj: # 空数组默认推断为string类型数组,可按需调整 item_schema = {"type": "string", "required": True} else: # 假设数组内元素类型一致,取第一个元素推断类型 item_schema = infer_schema(obj[0], required=True) return { "type": "array", "items": item_schema, "required": required } # 处理对象(AVRO Record)类型 if isinstance(obj, dict): schema["type"] = "record" schema["name"] = record_name schema["fields"] = {} for key, value in obj.items(): # 这里简化处理为所有字段必填,若要支持可选字段,可扩展多样本统计逻辑 schema["fields"][key] = infer_schema(value, f"{record_name}_{key}", required) return schema # 未知类型默认转为string return {"type": "string", "required": required} # 测试生成Schema if __name__ == "__main__": json_content = { "timestamp": 1661193367, "plant": "M14", "group": "g3", "device": "Sensor3", "parameter": "Power", "value": 100.3000000000001, "unit": "V", "limits": { "class": "A", "high": 200.2, "low": 32 }, "portfolio": { "id": "M00" } } avro_schema = infer_schema(json_content) write_schema(avro_schema, "avro_schema.json") print(json.dumps(avro_schema, indent=4))
Golang实现示例
以下是Golang版本的实现,通过递归反射推断JSON结构:
package main import ( "encoding/json" "os" "reflect" ) // AvroSchema 定义AVRO Schema的结构 type AvroSchema struct { Type string `json:"type"` Name string `json:"name,omitempty"` Fields map[string]AvroSchema `json:"fields,omitempty"` Items *AvroSchema `json:"items,omitempty"` Required bool `json:"required"` } func inferSchema(obj interface{}, recordName string, required bool) AvroSchema { v := reflect.ValueOf(obj) switch v.Kind() { case reflect.Int, reflect.Int8, reflect.Int16, reflect.Int32, reflect.Int64: return AvroSchema{Type: "long", Required: required} case reflect.Float32, reflect.Float64: return AvroSchema{Type: "double", Required: required} case reflect.String: return AvroSchema{Type: "string", Required: required} case reflect.Bool: return AvroSchema{Type: "boolean", Required: required} case reflect.Slice: var itemSchema AvroSchema if v.Len() > 0 { itemSchema = inferSchema(v.Index(0).Interface(), recordName+"_item", true) } else { // 空数组默认string类型元素 itemSchema = AvroSchema{Type: "string", Required: true} } return AvroSchema{Type: "array", Items: &itemSchema, Required: required} case reflect.Map: fields := make(map[string]AvroSchema) for _, key := range v.MapKeys() { keyStr := key.String() val := v.MapIndex(key).Interface() fields[keyStr] = inferSchema(val, recordName+"_"+keyStr, required) } return AvroSchema{Type: "record", Name: recordName, Fields: fields, Required: required} default: return AvroSchema{Type: "string", Required: required} } } func main() { jsonContent := map[string]interface{}{ "timestamp": 1661193367, "plant": "M14", "group": "g3", "device": "Sensor3", "parameter": "Power", "value": 100.3000000000001, "unit": "V", "limits": map[string]interface{}{ "class": "A", "high": 200.2, "low": 32, }, "portfolio": map[string]interface{}{ "id": "M00", }, } schema := inferSchema(jsonContent, "Record", true) file, err := os.Create("avro_schema_go.json") if err != nil { panic(err) } defer file.Close() encoder := json.NewEncoder(file) encoder.SetIndent("", " ") if err := encoder.Encode(schema); err != nil { panic(err) } }
内容的提问来源于stack exchange,提问作者Anand Narendra
相关产品推荐
相关产品推荐

