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

寻求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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 07:39:54