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

Pandas多字典填充值为列、内键为子列的DataFrame方法

Kafka持续OHLCV消息流的MultiIndex DataFrame填充方案

前置初始化代码(与问题描述规范完全一致)

import pandas as pd

# 两个交易所推送的OHLCV原始消息字典
dict1 = {
    "BTC/USDT": {"open": 42000.1, "high": 42300.5, "low": 41900.2, "close": 42250.8, "volume": 1256.3},
    "ETH/USDT": {"open": 2250.3, "high": 2275.9, "low": 2240.1, "close": 2268.4, "volume": 9872.6}
}
dict2 = {
    "BTC/USDT": {"open": 42001.2, "high": 42310.1, "low": 41895.7, "close": 42255.3, "volume": 1302.9},
    "ETH/USDT": {"open": 2251.0, "high": 2276.4, "low": 2239.2, "close": 2269.1, "volume": 9901.4}
}

# 构造一级为交易对标识、二级为OHLCV子字段的多级列索引
cols = pd.MultiIndex.from_product(
    [list(dict1.keys()), ["open", "high", "low", "close", "volume"]],
    names=["symbol", "ohlcv_field"]
)
# 初始化空DataFrame,行索引后续绑定Kafka消息时间戳/消费位点
df = pd.DataFrame(columns=cols)

可行填充方案

方案1:逐消息流式赋值(适配实时消费场景)

Kafka持续消费场景下无需等攒批,每消费到一条消息直接按列层级映射赋值即可,内存开销小,逻辑简单:

# 消费到dict1对应消息,以消息生产时间戳为行索引
msg1_ts = 1718000001000
for symbol, ohlcv_items in dict1.items():
    for field, val in ohlcv_items.items():
        df.loc[msg1_ts, (symbol, field)] = val

# 消费到dict2对应消息
msg2_ts = 1718000002000
for symbol, ohlcv_items in dict2.items():
    for field, val in ohlcv_items.items():
        df.loc[msg2_ts, (symbol, field)] = val

赋值完成后可统一执行df = df.astype(float)转换数值类型,避免逐格写入生成object类型影响后续计算。

方案2:批量转换拼接(适配小批量攒批消费场景)

如果消费端配置了攒批逻辑,可先将单条消息字典转换为匹配目标列结构的单行DataFrame,再批量拼接,执行效率高于逐单元格赋值:

def ohlcv_msg_to_df(msg: dict, row_idx) -> pd.DataFrame:
    row_data = []
    col_tuples = []
    # 严格按照初始化时的列顺序拼接数据
    for symbol in dict1.keys():
        for f in ["open", "high", "low", "close", "volume"]:
            col_tuples.append((symbol, f))
            row_data.append(msg[symbol][f])
    return pd.DataFrame(
        [row_data],
        index=[row_idx],
        columns=pd.MultiIndex.from_tuples(col_tuples, names=["symbol", "ohlcv_field"])
    )

# 批量拼接所有已攒的消息
batch_df_list = [
    ohlcv_msg_to_df(dict1, 1718000001000),
    ohlcv_msg_to_df(dict2, 1718000002000)
]
df = pd.concat([df] + batch_df_list)

注意事项

  • 列定位必须使用(一级列名, 二级列名)的元组格式,否则会生成冗余新列,破坏原有MultiIndex结构
  • 行索引优先选用Kafka消息自带的生产时间戳或消费位点,方便后续做时序对齐、消息去重、故障回溯
  • 若存在部分交易对字段缺失的情况,赋值前需做缺省值填充,避免触发KeyError

内容的提问来源于stack exchange,提问作者sci9

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 14:06:25