第六阶段 54 · ingest pipeline 摄取管道(写入时加工数据)
54 · ingest pipeline 摄取管道写入时加工数据阶段第六阶段 / 进阶专题ESingest pipeline processors | PostgreSQLBEFORE INSERT触发器 / ETL 转换1. 概念ingest pipeline让文档在写入索引之前先过一串processor处理器做加工补字段、改类型、拆分、脱敏、算派生值、调用 enrich 补维表……全部在 ES 侧完成应用端不用写这些转换逻辑。一句话写入时的轻量 ETL跑在 ES 协调节点上。常用 processorset、rename、convert、date、grok正则解析日志、split、gsub、scriptPainless、enrich补维表、remove。2. PostgreSQL 对照-- BEFORE INSERT 触发器写入前加工CREATEFUNCTIONfill_defaults()RETURNStriggerAS$$BEGINNEW.created_at :now();NEW.amount :NEW.price*NEW.qty;RETURNNEW;END;$$LANGUAGEplpgsql;CREATETRIGGERt BEFOREINSERTONsalesFOR EACH ROWEXECUTEFUNCTIONfill_defaults();ingest pipeline 就是 ES 版的「写入前触发器 / ETL 转换」但配置化、可复用。3. ES DSL3.1 定义 pipelinePUT _ingest/pipeline/sales_pipeline { description: 写入前补字段、算金额、打时间戳, processors: [ { set: { field: ingested_at, value: {{_ingest.timestamp}} } }, { convert: { field: qty, type: integer } }, { script: { source: ctx.amount ctx.price * ctx.qty } }, { remove: { field: tmp_raw, ignore_missing: true } } ], on_failure: [ // 出错兜底别丢数据 { set: { field: _ingest_error, value: {{ _ingest.on_failure_message }} } } ] }3.2 用 pipeline 写入# 单条写入时指定 POST sales_idx/_doc?pipelinesales_pipeline { price: 10, qty: 5, tmp_raw: x } # 写完得到 amount50、ingested_at... # 或把它设为索引默认之后所有写入自动走 PUT sales_idx/_settings { index.default_pipeline: sales_pipeline }3.3 用 _simulate 先试跑强烈推荐POST _ingest/pipeline/sales_pipeline/_simulate { docs: [ { _source: { price: 10, qty: 5 } } ] }4. Spring Boot 实现ComponentpublicclassDoc54IngestPipeline{AutowiredprivateElasticsearchClientelasticsearchClient;/** 创建/更新 pipeline定义外置成 JSONwithJson 直灌 */publicvoidputPipeline(Stringid,StringpipelineJson)throwsIOException{elasticsearchClient.ingest().putPipeline(p-p.id(id).withJson(newStringReader(pipelineJson)));}/** 写入时指定 pipeline也可在 bulk 的每个操作上带 pipeline */publicvoidindexWithPipeline(Stringindex,Stringpipeline,MapString,Objectdoc)throwsIOException{elasticsearchClient.index(i-i.index(index).pipeline(pipeline)// ← 走摄取管道加工.document(doc));}}ingest 客户端在elasticsearchClient.ingest()。bulk也支持在请求或每条操作上指定pipeline。5. 坑与最佳实践先_simulate再上线processor 顺序、字段名错很常见模拟能省大量返工。一定配on_failure默认某条处理失败会整条写入失败兜底把错误记下来别丢数据。default_pipeline方便但隐蔽设成索引默认后所有写入都会走排查问题记得想到它。grok解析日志很强但慢正则别写太贪婪结构化数据优先用dissect。ingest 跑在协调/ingest 节点重加工会占 CPU大流量考虑专用 ingest 节点。和 enrich 配合写入时补维表字段第 51 篇就是通过 ingest 的enrichprocessor。

相关新闻

Ubuntu 14.04解决Firefox无法播放小红书视频问题

Ubuntu 14.04解决Firefox无法播放小红书视频问题

1. 问题现象与背景分析最近在Ubuntu 14.04.6系统上使用自带的Firefox浏览器访问小红书时,发现视频内容无法正常播放。这个问题困扰了不少还在使用这个经典LTS版本的用户。作为一款2019年发布的系统,Ubuntu 14.04.6虽然稳定,但内置的软件包版本…

2026/8/4 11:43:07 阅读更多
木质也能做防火门?很多人都不知道

木质也能做防火门?很多人都不知道

多数人存在认知误区,认为防火门只能采用钢制材质,实际上符合国标要求的木质防火门早已广泛应用,依据 GB12955-2024,木质防火门属于正规被动防火构件,大量使用于酒店、写字楼、住宅精装区域。木质防火门并非普通实木门简…

2026/8/4 11:43:07 阅读更多
微信公众号数据采集:3步实现自动化运营分析

微信公众号数据采集:3步实现自动化运营分析

微信公众号数据采集:3步实现自动化运营分析 【免费下载链接】wechat_articles_spider 微信公众号文章的爬虫 项目地址: https://gitcode.com/gh_mirrors/we/wechat_articles_spider 你是否曾为手动统计公众号数据而烦恼?每天需要打开几十篇文章&a…

2026/8/4 12:43:10 阅读更多
3分钟搞定!QQ空间历史说说完整备份终极指南

3分钟搞定!QQ空间历史说说完整备份终极指南

3分钟搞定!QQ空间历史说说完整备份终极指南 【免费下载链接】GetQzonehistory 获取QQ空间发布的历史说说 项目地址: https://gitcode.com/GitHub_Trending/ge/GetQzonehistory 你是否曾想过,那些年发过的QQ空间说说,那些记录青春的文字…

2026/8/3 12:53:38 阅读更多
AMAT 0100-02186 I/O 分配 PCB

AMAT 0100-02186 I/O 分配 PCB

AMAT 0100-02186 I/O分配PCB板是应用材料(Applied Materials)公司生产的一款用于半导体设备的I/O信号分配电路板。该型号(0100-02186)的核心特点如下:专用于Endura等半导体工艺腔室。集成信号路由与分配功能。连接控制…

2026/8/3 19:34:52 阅读更多
Nissei Corp FFMN-32L-10-T0 40AX 三相异步电动机

Nissei Corp FFMN-32L-10-T0 40AX 三相异步电动机

Nissei Corp FFMN-32L-10-T0 40AX 三相异步电动机是日本日清(Nissei)品牌的一款工业用三相异步电机,适用于自动化设备及通用机械驱动。该型号(FFMN-32L-10-T0 40AX)的核心特点如下:三相交流异步电动机。额定…

2026/8/3 19:34:54 阅读更多