基于Apache NiFi高效处理百万级HTTP API请求并写入PostgreSQL
Apache NiFi 处理百万级HTTP请求+PostgreSQL插入的最优方案
核心架构思路
针对每日100万次请求、按时段分组、URL每日变更的场景,采用**"URL源管理-定时调度-请求分组-数据解析-批量入库"**的流水线模式,兼顾灵活性与性能。
详细步骤与处理器选型
1. URL源管理(应对每日200-300个URL变更)
- 用PostgreSQL存储URL清单:将URL、对应调度时段、剩余请求次数(初始设为5)存入数据库,方便每日更新;也可维护本地CSV文件,用
ListFile读取。 - 处理器:
ExecuteSQL(从数据库拉取每日有效URL,带时段标签)或ListFile+FetchFile(读取CSV),配合UpdateAttribute给每个URL添加时段标记(如morning/afternoon/evening)、剩余请求次数属性。
2. 定时调度与请求分组
- 按时段触发:用
Timer处理器设置多个调度规则(如早8点、午12点、晚6点),对应不同时段的URL分组。 - 分流控制:用
RouteOnAttribute根据URL的时段标记属性分流到对应请求分支,避免跨时段请求。 - 次数控制:每次请求后用
UpdateAttribute将剩余请求次数减1,当次数为0时路由到归档分支;用DistinctRecord确保当前时段内每个URL只处理一次。
3. 百万级HTTP请求优化
- 异步并发:用
InvokeHTTP开启异步模式,根据API限流设置合适并发数(建议200-500),避免压垮服务。 - 失败重试:添加
RetryFlowFile,对5xx、超时等失败请求自动重试,设置30秒间隔和3次最大重试次数。 - 流量整形:用
RateLimiter控制请求速率,匹配API限流规则,防止被封禁。
4. JSON数据解析与字段提取
- 处理器:
JoltTransformJSON,编写Jolt规范提取10个目标字段,示例规范:
{ "operation": "shift", "spec": { "field1": "db_field1", "field2": "db_field2", // 其余8个字段映射同理 } }
- 验证:用
ValidateRecord检查提取后的字段完整性,缺失字段的FlowFile路由到异常分支。
5. PostgreSQL批量入库
- 批量合并:用
MergeRecord将500-1000条解析后的FlowFile合并为批量数据,提升入库效率。 - 处理器:
PutDatabaseRecord,配置PostgreSQL连接池,选择JSONTreeReader作为记录读取器,映射解析字段到数据库表字段,开启批量插入模式,减少数据库连接开销。
6. 监控与异常处理
- 监控:用
MonitorActivity跟踪请求、入库成功率,配合NiFi自带面板查看吞吐量。 - 异常处理:将请求失败、解析失败的FlowFile路由到
PutFile存储到异常目录,方便后续排查。
关键模式说明
- 增量更新模式:每日用
ExecuteSQL拉取新增/变更URL,无需全量重载,提升效率。 - 分段调度模式:多
Timer+RouteOnAttribute实现按时段分组请求,避免集中请求触发API限流。 - 批量处理模式:从解析到入库全流程批量操作,降低NiFi内部调度和数据库压力。
内容的提问来源于stack exchange,提问作者özüm
相关产品推荐
相关产品推荐

