使用Python向Elasticsearch发送JSON并将t字段设置为时间戳索引
问题根因
- 未提前为Elasticsearch索引定义显式映射,ES自动识别13位数值格式的
t字段为Long类型,不会默认转换为日期格式 - 原有写入逻辑是将完整的API响应作为单个文档存入ES,
results数组内的多条K线数据被嵌套在同一个文档中,无法作为独立的时间序列节点被查询 - 代码中使用的
doc_type参数在ES 7.x及以上版本已经废弃,无需指定
修复步骤
1. 删除原有错误索引
首先删除已经生成错误映射的core索引,执行以下请求:
DELETE /core
你可以通过curl、Kibana Dev Tools执行该请求,也可以在Python代码中调用es.indices.delete(index='core', ignore=[400, 404])实现。
2. 创建带正确映射的新索引
提前定义索引映射,显式指定t字段为日期类型,格式匹配13位毫秒级时间戳:
PUT /core { "mappings": { "properties": { "ticker": { "type": "keyword" }, "adjusted": { "type": "boolean" }, "v": { "type": "long" }, "vw": { "type": "float" }, "a": { "type": "float" }, "o": { "type": "float" }, "c": { "type": "float" }, "h": { "type": "float" }, "l": { "type": "float" }, "t": { "type": "date", "format": "epoch_millis" }, "n": { "type": "long" } } } }
3. 修改Collector.py写入逻辑
原有代码存在包导入位置错误、写入逻辑不合理的问题,修改后代码如下:
# 包导入统一放在函数顶部,避免重复导入 import json import os import requests from elasticsearch import Elasticsearch from flask import render_template def collect(): key = "" url = "https://api" payload={} headers = {} response = requests.request("GET", url, headers=headers, data=payload) j = response.json() data = response.text # 保存原始JSON到本地 with open('data.json', 'w', encoding='utf-8') as f: json.dump(j, f, ensure_ascii=False, indent=4) # 初始化ES客户端 es = Elasticsearch([{'host': 'localhost', 'port': 9200}]) # 提取公共字段 common_fields = { "ticker": j["ticker"], "adjusted": j["adjusted"] } doc_id = 1 # 遍历results数组,每条K线作为独立文档写入ES for kline in j["results"]: # 合并公共字段和K线字段 doc = {**common_fields, **kline} es.index(index='core', id=doc_id, body=doc) doc_id += 1 return render_template('show.html', data=data)
修改说明:
- 移除了已废弃的
doc_type参数 - 不再将完整JSON作为单文档写入,拆分
results数组的每条K线为独立文档 - 每条K线文档自动拼接股票代码、是否复权等公共属性,方便后续多股票数据查询
功能验证
写入完成后可执行GET /core/_mapping查询索引映射,确认t字段类型为date,之后可直接使用时间范围语法查询数据,实现时间序列相关分析功能。
内容的提问来源于stack exchange,提问作者Luke Patrick
相关产品推荐
相关产品推荐

