elasticsearch-py使用async_bulk报错:actions中缺少index字段
问题原因
原生AsyncElasticsearch.bulk()接收的是符合ES原生bulk API要求的NDJSON格式字符串序列,需要手动拆分操作元数据行和文档数据行;但async_bulk作为上层封装工具,会自动完成序列化、格式拼接逻辑,要求传入的action为Python字典类型,而非JSON字符串,这是你报错的核心原因。
action数据结构定义
async_bulk要求每个action是Python字典,结构分为两类字段:
- 操作元数据字段(带下划线前缀):
_index:目标索引名_op_type:操作类型,可选值为index/create/update/delete,默认值为index,无需额外指定_id:文档ID,可选,不指定则ES自动生成
- 文档数据字段:可直接与元数据字段平级放置,也可统一放在
_source字段中
修正方案
方案1:在action中单独指定索引
def _rec_to_actions(self, df): for record in df.to_dict(orient="records"): yield { "_index": self.index, **record } async def send_to_elasticsearch(self, df: DataFrame): await async_bulk(self.elastic_client, self._rec_to_actions(df))
方案2:全局指定索引(更简洁)
如果所有数据都写入同一个索引,可以直接在调用async_bulk时传入公共index参数,不需要在每个action中重复指定:
def _rec_to_actions(self, df): # 直接返回文档内容字典即可 yield from df.to_dict(orient="records") async def send_to_elasticsearch(self, df: DataFrame): await async_bulk( self.elastic_client, self._rec_to_actions(df), index=self.index )
内容的提问来源于stack exchange,提问作者zar3bski
相关产品推荐
相关产品推荐

