如何转换DynamoDB数据结构适配Elasticsearch写入查询需求
DynamoDB 流数据转 Elasticsearch 目标格式实现方案
你拿到的是DynamoDB Streams输出的标准带类型标记的记录格式,处理逻辑拆成3个固定步骤即可,不需要额外依赖复杂组件:
- 提取元数据:把最外层
index对象内的_index、_type、_id三个ES索引元字段直接提取到结果根层级 - 解析主键字段:遍历
Keys下的所有字段,去掉DynamoDB的类型包裹层(比如{"S": "sk001"}里的S是类型标记,代表字符串类型,实际值是sk001),把键值对直接放到结果根层级 - 解析业务字段:遍历
NewImage下的所有业务字段,同样去掉类型包裹层提取实际值,合并到结果根层级;不需要的字段(比如示例里的ApproximateCreationDateTime、price)可以直接过滤丢弃
注:你提供的原始JSON和目标JSON都存在语法错误,
"_type": "_doc后缺少闭合引号和逗号,实际处理前需要先确保输入JSON是合法可解析的结构。
以下是可直接运行的Python实现代码:
def parse_dynamodb_typed_item(typed_item: dict) -> dict: """剥离DynamoDB字段的类型包装,返回原生键值对""" plain_item = {} for field, type_wrapper in typed_item.items(): # 按需扩展支持的DynamoDB类型即可,覆盖常用场景 if "S" in type_wrapper: plain_item[field] = type_wrapper["S"] elif "N" in type_wrapper: # 数字类型可根据业务需要转int/float plain_item[field] = int(type_wrapper["N"]) elif "BOOL" in type_wrapper: plain_item[field] = type_wrapper["BOOL"] return plain_item def transform(raw_ddb_record: dict, exclude_fields: list = None) -> dict: exclude_fields = exclude_fields or [] es_doc = {} # 提取ES索引元信息 es_doc.update(raw_ddb_record["index"]) # 提取主键字段 es_doc.update(parse_dynamodb_typed_item(raw_ddb_record["Keys"])) # 提取业务字段,过滤不需要的字段 parsed_new_image = parse_dynamodb_typed_item(raw_ddb_record["NewImage"]) for field in exclude_fields: parsed_new_image.pop(field, None) es_doc.update(parsed_new_image) return es_doc # 调用示例 if __name__ == "__main__": raw_record = { "index": { "_index": "data-index", "_type": "_doc", "_id": "verBionoVub" }, "ApproximateCreationDateTime": 1645005038, "Keys": { "sk": {"S": "sk001"}, "pk": {"S": "pk001"} }, "NewImage": { "res_type": {"S": "message"}, "author": {"S": "user"}, "price": {"N": "15"}, "product_id": {"S": "B016JOMAEE"} } } # 不需要price字段就传入exclude_fields过滤 result = transform(raw_record, exclude_fields=["price"]) print(result)
运行后输出的结果和你要求的目标格式完全一致,如果是用其他语言(比如Java、Go)处理,逻辑完全相同,只需要按照上述三个步骤实现类型剥离和字段合并即可。
内容的提问来源于stack exchange,提问作者vishnuvid
相关产品推荐
相关产品推荐

