Flume HTTPSource 与 HTTP Sink 实践:构建实时数据接收网关与推送端点

发布时间:2026/8/31 1:02:29

Flume HTTPSource 与 HTTP Sink 实践:构建实时数据接收网关与推送端点
Flume HTTPSource 与 HTTP Sink 实践构建实时数据接收网关与推送端点Flume HTTPSource 与 HTTP Sink 概述Apache Flume 是一个分布式、可靠、可扩展的服务用于高效地收集、聚合和移动大量日志数据。在实时数据处理场景中Flume 的 HTTPSource 和 HTTP Sink 组件提供了通过 HTTP 协议进行数据接收和推送的能力。HTTPSource 允许 Flume 接收来自外部 HTTP 请求的数据适用于将 Web 应用、移动应用等产生的日志实时接入数据管道。HTTP Sink 则使 Flume 能够将处理后的数据通过 HTTP 协议发送到外部服务如 Elasticsearch、Kafka 或其他自定义 API 端点。这两种组件的结合使用可以构建灵活的数据处理网关实现数据的实时采集、转换和分发满足现代分布式系统中对实时数据流处理的需求。HTTPSource 实践构建实时数据接收网关HTTPSource 是 Flume 的一个内置 Source 组件通过 HTTP 协议接收数据。配置和使用 HTTPSource 接收 HTTP 请求需要以下步骤a. 在 Flume 配置文件中定义 HTTPSourceproperties# 定义源a1.sources r1a1.sources.r1.type org.apache.flume.source.http.HTTPSourcea1.sources.r1.bind 0.0.0.0a1.sources.r1.port 8080a1.sources.r1.handler org.apache.flume.source.http.JSONEventServleta1.sources.r1.handler.type jsona1.sources.r1.channels c1以上配置创建了一个监听在 0.0.0.0:8080 的 HTTPSource使用 JSONEventServlet 处理请求并将数据发送到通道 c1。b. 启动 Flume 代理bashflume-ng agent --conf ./conf --conf-file ./http-source.conf --name a1 -Dflume.root.loggerINFO,consolec. 使用 curl 或其他 HTTP 客户端发送数据bashcurl -X POST -H Content-Type: application/json -d {timestamp:2023-05-01T12:00:00, event:user_login, user:testuser} http://localhost:8080d. 验证数据是否被接收和处理配置一个 Memory Channel 和 Logger Sink 来验证数据流properties# 定义通道a1.channels c1a1.channels.c1.type memorya1.channels.c1.capacity 1000a1.channels.c1.transactionCapacity 100# 定义接收器a1.sinks k1a1.sinks.k1.type loggera1.sinks.k1.channel c1通过以上配置HTTPSource 接收到的数据将被发送到 Memory Channel最终通过 Logger Sink 输出到控制台。在实际应用中可以将 Logger Sink 替换为 HDFS、Kafka 或其他 Sink将数据持久化或进一步处理。HTTP Sink 实践构建实时数据推送端点HTTP Sink 是 Flume 的一个内置 Sink 组件通过 HTTP 协议发送数据到外部服务。配置和使用 HTTP Sink 需要以下步骤a. 在 Flume 配置文件中定义 HTTPSinkproperties# 定义源a1.sources r1a1.sources.r1.type execa1.sources.r1.command tail -F /var/log/flume/test.loga1.sources.r1.channels c1# 定义通道a1.channels c1a1.channels.c1.type memorya1.channels.c1.capacity 1000a1.channels.c1.transactionCapacity 100# 定义接收器a1.sinks k1a1.sinks.k1.type org.apache.flume.sink.http.HttpSinka1.sinks.k1.channel c1a1.sinks.k1.httpEndpoint http://localhost:8081/eventsa1.sinks.k1.httpMethod POSTa1.sinks.k1.contentType application/jsona1.sinks.k1.connectTimeout 30000a1.sinks.k1.requestTimeout 30000a1.sinks.k1.connectRetryDelay 10000a1.sinks.k1.defaultBackoff truea1.sinks.k1.maxBackoff 10000a1.sinks.k1.serializer org.apache.flume.sink.http.HttpServletRequestSerializer以上配置创建了一个 HTTPSink将数据通过 POST 请求发送到 http://localhost:8081/events使用 JSON 格式。b. 启动 Flume 代理bashflume-ng agent --conf ./conf --conf-file ./http-sink.conf --name a1 -Dflume.root.loggerINFO,consolec. 创建一个简单的 HTTP 服务来接收数据使用 Node.js 创建一个简单的 HTTP 服务javascriptconst http require(http);const server http.createServer((req, res) {if (req.method POST req.url /events) {let body ;req.on(data, chunk {body chunk.toString();});req.on(end, () {console.log(Received data:, body);res.writeHead(200);res.end(OK);});} else {res.writeHead(404);res.end(Not Found);}});server.listen(8081, () {console.log(Server running at http://localhost:8081/);});d. 验证数据是否被发送和接收向 /var/log/flume/test.log 文件中添加内容观察 Flume 是否将数据发送到 HTTP 服务以及 HTTP 服务是否接收到数据。完整实例构建实时数据流处理系统结合前面的 HTTPSource 和 HTTP Sink我们可以构建一个完整的实时数据流处理系统该系统接收来自 Web 应用的日志数据经过处理后将数据发送到 Elasticsearch 进行存储和分析。a. 配置 Flume 代理properties# 定义源a1.sources r1a1.sources.r1.type org.apache.flume.source.http.HTTPSourcea1.sources.r1.bind 0.0.0.0a1.sources.r1.port 8080a1.sources.r1.handler org.apache.flume.source.http.JSONEventServleta1.sources.r1.handler.type jsona1.sources.r1.channels c1# 定义通道a1.channels c1a1.channels.c1.type memorya1.channels.c1.capacity 1000a1.channels.c1.transactionCapacity 100# 定义接收器a1.sinks k1a1.sinks.k1.type org.apache.flume.sink.http.HttpSinka1.sinks.k1.channel c1a1.sinks.k1.httpEndpoint http://elasticsearch:9200/logs/_doca1.sinks.k1.httpMethod POSTa1.sinks.k1.contentType application/jsona1.sinks.k1.connectTimeout 30000a1.sinks.k1.requestTimeout 30000a1.sinks.k1.connectRetryDelay 10000a1.sinks.k1.defaultBackoff truea1.sinks.k1.maxBackoff 10000a1.sinks.k1.serializer org.apache.flume.sink.http.HttpRequestBodySerializerb. 启动 Flume 代理bashflume-ng agent --conf ./conf --conf-file ./flume.conf --name a1 -Dflume.root.loggerINFO,consolec. 使用 curl 发送数据bashcurl -X POST -H Content-Type: application/json -d {timestamp: 2023-05-01T12:00:00,level: INFO,message: User login,user: testuser,ip: 192.168.1.100} http://localhost:8080d. 验证数据是否被存储到 Elasticsearch使用 Elasticsearch 的 REST API 或 Kibana 检查数据是否被正确存储bashcurl -X GET http://elasticsearch:9200/logs/_search?pretty注意事项与最佳实践在使用 Flume 的 HTTPSource 和 HTTP Sink 时需要注意以下几点a.性能优化合理配置通道容量和事务大小避免数据丢失或性能瓶颈对于高并发场景考虑使用多通道或多个 Flume 代理实例b.错误处理配置适当的重试机制和超时设置实现监控和告警机制及时发现和处理数据流异常c.安全考虑对 HTTPSource 启用 HTTPS 和基本认证对敏感数据进行加密处理d.数据格式统一数据格式便于后续处理和分析考虑使用 Schema Registry 管理数据结构变更e.扩展性使用 Load Balance Channel 或 Fanout Channel 实现数据分流考虑使用 Flume NG 集群部署提高可靠性最小示例与注意事项HTTPSource 配置文件 (http-source.conf):# 定义源 a1.sources r1 a1.sources.r1.type org.apache.flume.source.http.HTTPSource a1.sources.r1.bind 0.0.0.0 a1.sources.r1.port 8080 a1.sources.r1.handler org.apache.flume.source.http.JSONEventServlet a1.sources.r1.handler.type json a1.sources.r1.channels c1 # 定义通道 a1.channels c1 a1.channels.c1.type memory a1.channels.c1.capacity 1000 a1.channels.c1.transactionCapacity 100 # 定义接收器 a1.sinks k1 a1.sinks.k1.type logger a1.sinks.k1.channel c1启动命令:flume-ng agent --conf ./conf --conf-file ./http-source.conf --name a1 -Dflume.root.loggerINFO,console发送数据:curl -X POST -H Content-Type: application/json -d {event:test} http://localhost:8080注意事项:确保防火墙开放了 Flume 监听的端口检查 Flume 版本HTTPSource 和 HTTP Sink 的类名可能随版本变化对于生产环境应考虑配置多个通道和备份接收器以提高可靠性监控 Flume 的内存使用情况避免内存溢出大数据量场景下考虑增加 batch-size 参数提高吞吐量数据流程图:POST请求接收事件传输数据HTTP请求HTTP客户端HTTPSourceChannelHTTPSink外部服务

相关新闻

Delphi数据库结构同步利器:Clever Database Comparer控件详解与应用

Delphi数据库结构同步利器:Clever Database Comparer控件详解与应用

2026/8/31 0:52:28

简介:这是一套专为Delphi开发者设计的数据库结构比对与同步工具——Clever Database Comparer v8.2.955.0,兼容Delphi 7至最新版13.1(Florence),面向中高级数据库应用开发人员及DBA,解决多环境数据库Schema…

【原创】基于AI大模型+SpringBoot+Vue的电影院选座购票网站(设计与实现)

【原创】基于AI大模型+SpringBoot+Vue的电影院选座购票网站(设计与实现)

2026/8/31 0:52:28

摘要:随着行业信息化建设持续推进,电影院选座购票网站相关业务对线上协同与数据沉淀的要求不断提高。传统线下或分散式办理方式存在流程繁琐、信息滞后、协作成本高、过程难追溯等弊端,难以适应便捷化、可管理的业务服务需求。同类课题亦多见…

【原创】基于微信小程序+AI大模型+uni-app的电影院选座购票小程序(设计与实现)

【原创】基于微信小程序+AI大模型+uni-app的电影院选座购票小程序(设计与实现)

2026/8/31 0:52:28

摘要:随着行业信息化建设持续推进,电影院选座购票网站相关业务对线上协同与数据沉淀的要求不断提高。传统线下或分散式办理方式存在流程繁琐、信息滞后、协作成本高、过程难追溯等弊端,难以适应便捷化、可管理的业务服务需求。同类课题亦多见…

OpenAI 断供 Cursor:AI 编程的模型供应链断了

OpenAI 断供 Cursor:AI 编程的模型供应链断了

2026/8/31 2:12:32

8 月 28 日,OpenAI 通知 SpaceX,计划终止向 AI 编程工具 Cursor 提供 OpenAI 模型的合同,拟定终止日期是 2026 年 11 月 12 日。消息一出,天天用 Cursor 写代码的人多少会愣一下:我天天用的工具,底层模型还…

STM32结合RFID图书管理系统:从硬件选型到云端联调全解析

STM32结合RFID图书管理系统:从硬件选型到云端联调全解析

2026/8/31 2:12:32

简介:本资源是一套基于STM32平台的物联网图书管理系统毕业设计实战案例,面向高校电子、通信、自动化及物联网相关专业本科生,解决图书馆场景下图书借还、身份识别与数据管理等核心问题,适用于毕业设计选题、课程设计实践及嵌入式开…

Python 批量图像处理实战:为活动返图打造肤色提亮与日系滤镜流水线

Python 批量图像处理实战:为活动返图打造肤色提亮与日系滤镜流水线

2026/8/31 2:12:32

最近“异环日本线下活动”的一组返图在社交平台上的讨论度很高,尤其是菌烨小姐姐还原的“真红”,服装细节、妆面质感以及神态都相当到位。我们在欣赏这类高质量返图时,如果切换回开发视角,会发现线下活动返图其实是一个非常典型的…

从虎扑评分看电竞社区数据产品:NIP vs WBG的赛后数据拆解

从虎扑评分看电竞社区数据产品:NIP vs WBG的赛后数据拆解

2026/8/31 2:12:32

如果只看比分,你会觉得这只是一场普通的 BO3 常规赛:NIP 2-1 WBG,三局打满,赢家带走胜利,输家回去复盘。但如果你把视线移到赛场之外的虎扑评分区,会发现这场比赛的热度远远超出“2-1”这个数字本身。选手评…

红外弱小目标检测与跟踪的Matlab实现:原理、代码与调参指南

红外弱小目标检测与跟踪的Matlab实现:原理、代码与调参指南

2026/8/31 2:12:32

简介:本资源面向图像处理初学者与红外目标跟踪研究者,提供一套完整、可直接运行的弱小目标检测与跟踪MATLAB实现方案,聚焦于低信噪比红外图像中的目标识别与运动轨迹估计问题。压缩包共7个文件,含3个核心M函数(主程序m…

AI模型安全扫描器:为何F1不如覆盖率与故障恢复重要

AI模型安全扫描器:为何F1不如覆盖率与故障恢复重要

2026/8/31 2:02:32

如果只用一个指标去衡量一款 AI 模型安全扫描器,你会选什么?我见过很多团队直接看 F1。理由很直接:F1 同时包含精确率和召回率,能用单一分数说明检测能力。但这个习惯放到 AI 模型安全扫描器上,往往会在生产环境里埋雷…

备战数据库管理工程师校招:索引、事务、备份恢复核心考点解析

备战数据库管理工程师校招:索引、事务、备份恢复核心考点解析

2026/8/31 1:38:25

每年校招季我都会接触不少准备数据库方向笔试的同学,看到最多的状态就是:简历上写着“熟悉 MySQL”“了解索引优化”,一碰到数据库管理工程师的笔试卷,却在索引、事务、锁、备份恢复这些题目上翻车。网易这套 2018 校园招聘数据库…

数字电路时序基石:深入理解建立时间与保持时间

数字电路时序基石:深入理解建立时间与保持时间

2026/8/30 0:01:07

1. 这不是“背公式”的事:时间参数到底在约束什么你翻过数字电路教材,一定见过这两个词:建立时间(Setup Time)和保持时间(Hold Time)。它们常被并列写在触发器(Flip-Flop&#xff09…

蓝桥杯国赛超声波测距机:从单片机原理到嵌入式系统实战

蓝桥杯国赛超声波测距机:从单片机原理到嵌入式系统实战

2026/8/30 0:01:07

1. 项目缘起:从赛题到超声波测距机的诞生第八届蓝桥杯单片机设计与开发国赛的题目,我至今记忆犹新。它没有直接给出一个花哨的名字,而是用“超声波测距机”这个朴实无华的功能描述,精准地勾勒出了考核的核心。对于当时备赛的我而言…

MCU无DAC如何用定时器+DMA 2D输出高保真任意波形

MCU无DAC如何用定时器+DMA 2D输出高保真任意波形

2026/8/31 0:02:27

接到一个仪表类项目,要在 LAT1189 上输出几种不同波形:正弦、三角、带可调死区的脉冲,频率和幅度都得能实时改。板子上没有 DAC,就一个定时器加几个 DMA 通道。我一开始觉得在定时器中断里改比较寄存器也能应付,后来把…

Cortex-M3 Flash下载失败?从编程错误标志到供电瞬态排查

Cortex-M3 Flash下载失败?从编程错误标志到供电瞬态排查

2026/8/31 0:02:27

前两周调试一块带着Cortex-M3内核的板子,IDE里下载固件时突然弹出一行刺眼的错误: error: flash download failed - cortex-m3 。这种报错在嵌入式开发里太常见了,常见到很多人第一反应就是换根数据线、重插一下调试器,但重启三…

STM32 TouchGFX屏幕切换Transition优化:原理、配置与排障实战

STM32 TouchGFX屏幕切换Transition优化:原理、配置与排障实战

2026/8/31 0:02:27

做STM32 GUI开发的朋友应该都有体会——界面搭得再漂亮,一旦屏幕切换卡成PPT,整个产品的档次瞬间就没了。早期我在LAT1212这个基于STM32的GUI工程上用TouchGFX做二次开发,最头疼的不是画界面,而是怎么让切换动画既流畅又自然。Tou…

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

摆脱论文困扰!盘点2026年全网爆红的的AI论文写作工具

2026/8/28 7:35:26

一天写完毕业论文在2026年已不再是天方夜谭。2026年最炸裂、实测能大幅提速的AI论文写作工具,覆盖选题构思、文献整理、内容生成、格式排版等核心场景,真正帮你高效搞定论文难题。 一、全流程王者:一站式搞定论文全链路(一天定稿首…

导师推荐!2026最新AI论文工具测评与实用推荐

导师推荐!2026最新AI论文工具测评与实用推荐

2026/8/28 7:34:51

2026年真正好用的AI论文工具,核心看生成的论文质量、低AI味、格式正确、学术适配四大指标。综合实测,千笔AI、ThouPen、豆包、DeepSeek、Grammarly 是当前最值得推荐的梯队,覆盖从免费到付费、从中文到英文、从文科到理工的全场景需求。 一、…

告别游戏崩溃:XCOM 2模组管理器的智能革命

告别游戏崩溃:XCOM 2模组管理器的智能革命

2026/8/28 7:34:35

告别游戏崩溃:XCOM 2模组管理器的智能革命 【免费下载链接】xcom2-launcher The Alternative Mod Launcher (AML) is a replacement for the default game launchers from XCOM 2 and XCOM Chimera Squad. 项目地址: https://gitcode.com/gh_mirrors/xc/xcom2-lau…