做大数据采集选型,大多数人开口第一句是“用 Flume 还是 Logstash?”,第二句是“要不要上 Kafka?”。我这些年接触过不少数据团队,发现这个讨论顺序本身就是错的——工具排名永远在需求之前,结果常常是拍脑袋定了一个框架,上线之后被各种边界情况反复吊打。这篇文章想把“大数据采集方案怎么选”这件事讲透:不是给你罗列“某某工具多牛”,而是给出一套能照着做、能打分的决策路径。
后面会按四条线展开:先讲选型前必须想清楚的三个前置问题,再拆解日志采集、数据库同步、消息缓冲三类主流工具的真实边界,然后给出一张可以直接套用的量化对比表,最后聊一聊方案上线后最容易翻车的四个环节。适合正在搭数仓、做数据中台,或者负责实时计算链路的数据工程师、架构师、技术负责人阅读。如果你正卡在“别人说什么好就用什么”的阶段,这篇大概率能帮你少走半年弯路。
1. 选型先问三个问题:数据从哪来、要送多快、落到哪去
很多选型会开成“工具推荐会”,这个习惯不太对。工具是需求的倒影,需求没理清,再强的框架也白搭。我习惯在动工具之前,强制自己和业务方把三个问题聊透:数据源到底是什么形态、数据从产生到可用能容忍多久、到了目标端之后以什么形式被消费。这三个问题直接决定你该往哪个技术栈靠,而不是反过来被某个框架绑定。下面逐个展开。
1.1 数据源类型决定了技术栈走向
数据源形态基本决定了采集器的长相,也是最不该偷懒的盘点环节。最常见的四类:
- 文本日志类。服务器日志、业务打印日志、网站埋点日志。这类数据本质是“追着文件跑”,需要 agent 驻留在业务机器上,或者用轮询扫目录。典型工具是 Filebeat、Flume、Logstash。选型的核心矛盾,是 agent 的资源占用和吞吐量之间的平衡,后面会细说。
- 数据库变更类。业务系统存在 MySQL、PostgreSQL、Oracle 里,你希望把新增、更新、删除的记录同步到数仓。这里要分成两条路:全量同步用 DataX、Sqoop 这类批式框架;增量变更用 Canal、Flink CDC 这类监听 binlog 的组件。很多团队以为“数据库同步”是一个工具能搞定的,实际上这两条路线常常要配合使用。
- 消息队列 / API 推送类。上游系统已经把数据打进 Kafka、RocketMQ,或者通过 HTTP 接口推给你。这种情况下“采集”其实退化成了消费和落盘,更多要考虑的是消费速度、幂等和背压,本身不需要重型 agent。
- IoT / 边缘设备类。车载、智能硬件、工业传感器,协议五花八门,常见 MQTT、Modbus、OPC UA。通常不能在设备上直接装 Flume,需要边缘网关先做协议接入,再汇总到流处理链路,和数据中心里采集日志完全两个物种。
我见过最典型的错配,是把 Flume agent 塞进 IoT 设备网关里,理由是“它支持 Kafka sink”。先不说 Java 进程的体量,单是设备掉线、弱网重传、流量计费这些问题,Flume 默认机制就顶不住。设备侧通常要先经 MQTT Broker,再做协议转换,这跟“采集服务端日志”根本不是一回事,选型时别混为一谈。
1.2 时效性要求:批、微批、还是秒级实时
第二个问题是“数据从产生到能用,业务能等多久”。这个问题不解决,很容易做出杀鸡用牛刀或者牛刀用成杀鸡的决策。
- T+1 离线。昨天夜里跑批,今天早上能看到报表。这种情况用 DataX/Sqoop 在低峰期做全量或增量同步最省心,省掉的不仅是实时计算的复杂度,还有一大笔故障恢复成本。
- 分钟级准实时。报表每小时刷新、风控策略每 5 分钟跑一次。可以用 Filebeat/Flume 采集 + Kafka 缓冲 + Spark Streaming 或 Flink 做微批。这个层级是大多数互联网业务的舒适区,投入产出比最高。
- 秒级实时。业务要求在数据产生后的几秒内就参与计算,比如实时反欺诈、实时大屏。这时候才需要 Flink CDC、Kafka Streams 这类真正的流式链路,并且要接受 checkpoint、状态后端、反压这些概念带来的运维负担。
有段时间“实时数仓”被讲得玄乎,有些团队为了在汇报里写上“实时”两个字,把所有链路都改成秒级流处理,结果吞吐上不去,人员也跟不上。我的观点是:实时能力是为具体业务场景买单的,不是为形容词买单的。如果业务方说“多等五分钟没问题”,那就老老实实做微批,省下的成本能干很多别的事。
1.3 目标端与数据形态:不是扔进HDFS就收工
采集数据的落点,直接影响工具选型。同样是“同步一张 MySQL 表”,落到不同目标端,最优解完全不同:
- 落到HDFS/Hive做离线数仓,DataX 天然支持,能直接写 parquet,还能控制分区;
- 落到Kafka给实时计算引擎消费,Flink CDC、Canal 反而是更顺手的入口;
- 落到ClickHouse、Doris做即席分析,要重点看采集工具的写入插件质量,DataX 对这两家都有专门支持,但有些工具只支持 JDBC 通用写入,性能差一个数量级;
- 去数据湖(Iceberg/Hudi)时,还要关心 CDC 产生的变更数据能不能转成 upsert 语义,这已经超出采集工具的职能,需要下游表格式配合。
很多人以为采集就是“把文件挪到一个地方”,忽略了目标端的数据格式。如果下游是 Hive 数仓,文本格式还是得转 parquet/ORC,这个转换放在采集阶段还是下游计算阶段,差别很大。建议在选型阶段就把目标端的 schema 和格式定下来,宁可初始多用点时间,也别等数据跑起来再返工。
2. 主流采集工具边界拆解:日志、库表、消息三类工具各管一段
理清前置问题之后,再看工具清单就不会晕了。我把常见工具按职责分成三类,每一类解决一个层面的问题,别指望一个工具包打天下。很多人天天纠结“Flume 和 Logstash 哪个更强”,其实那是把两个不同侧重点的工具强行拉到同一水平线比。真正的工程做法是让它们各管一段,通过组合来覆盖整条链路。
2.1 日志采集三兄弟:Filebeat、Logstash、Flume
这三者都做日志采集,但定位差异很明显。
| 工具 | 语言 | 资源占用 | 核心优势 | 弱势 |
|---|---|---|---|---|
| Filebeat | Go | 低 | 轻量级 agent,适合装在每台业务机器上;tail 文件、断点续读很稳 | 几乎不具备数据处理能力,最多做简单过滤 |
| Logstash | Java | 高 | 正则、Grok、字段转换能力强,插件生态丰富 | 吞吐有限,单实例性能一般;吃内存 |
| Flume | Java | 中高 | 多跳路由、Source/Channel/Sink 组件化,支持 Kafka 与 HDFS sink | agent 本体偏重,配置维护成本高,社区迭代慢 |
我的建议更偏向“组合拳”:业务机器上尽量只装 Filebeat,它用 Go 写的,在每台机器上跑大概只占二三十 MB 内存,它的工作就是把日志稳定地送出去;复杂的解析、清洗放到 Logstash 或者下游 Flink 做。早期很多团队习惯在所有节点装 Flume agent,一行配置改错,几百台机器要重启,维护成本让人崩溃。Filebeat + Kafka 再把日志交接到 Logstash/Flink,是目前更清晰的分工方式。
为什么不是 Logstash 直接作为采集 agent?因为它是 Java 进程,Grok 解析遇到高吞吐日志时 CPU 会明显走高,而且进程一挂,本地日志堆积的补偿机制相对笨重。作为集中式处理层它很有价值,但铺到每一台业务机,性价比不高。反过来,Filebeat 不适合做复杂解析,如果不想引入第二个处理层,小规模场景里让 Logstash 一台机器扛住几百 MB/s 以内的日志量,也不是不能接受。
2.2 数据库同步四大件:Sqoop、DataX、Canal、Flink CDC
数据库同步常被分成“全量 + 增量”两步走,工具选型也一样。先给结论,再解释为什么。
- Sqoop。基于 MapReduce 的老牌工具,现在社区基本不活跃,性能也不算好。除非你还维护着很老的 Hadoop 发行版,否则新项目不建议再选,踩坑成本比收益高。
- DataX。阿里开源的批式同步框架,代码简洁,插件化做得好,单机多线程就能跑出不错的速度。它对异构数据源(MySQL、Oracle、HDFS、Hive、ClickHouse 等)都有现成插件,日常全量同步和数据搬迁我用它最多。
- Canal。定位是 MySQL binlog 解析组件,把 binlog 变更解析后投递到 Kafka/RocketMQ 或者直接给下游。它是独立部署的 Java 服务,要自己维护集群、监听位点。
- Flink CDC。这个其实不是独立工具,而是 Flink 的源连接器,可以基于 binlog 直接把变更流拉进 Flink 做实时计算或写入数仓。相比 Canal,它天然具备 Flink 的 checkpoint、exactly-once、状态管理能力,链路更短,但对团队的 Flink 能力有要求。
| 工具 | 类型 | 全量 | 增量 | 实时 | 典型场景 |
|---|---|---|---|---|---|
| Sqoop | 批 | 支持 | 支持 | 否 | 遗留 Hadoop 生态 |
| DataX | 批 | 支持 | 支持 | 否 | 离线批量同步、异构搬迁 |
| Canal | 流 | 否 | 支持 | 是 | 监听 MySQL binlog,投递到 MQ |
| Flink CDC | 流 | 支持 | 支持 | 是 | 整库迁移、实时入湖入仓 |
有一个点被很多人忽略:DataX 的“增量”通常是按时间戳或 ID 轮询,这会持续给业务库加查询压力;而 binlog 方案对源库的侵入小得多。如果增量粒度很细、源库又特别忙,优先考虑 CDC 路线。反过来,如果你的源表数据经常被删除修改,CDC 能拿到完整历史变更,而轮询方式对“删除”是无能为力的。全量与增量之间的衔接,我在后面第四部分会专门讲,这里先记住“两条腿走路”这个原则。
2.3 Kafka 的角色:是缓冲通道,不是采集终点
一提到高吞吐采集,很多人第一反应就是“上 Kafka”。这个判断大方向对,但对“Kafka 在链路里到底该承担什么角色”想得往往不够清楚。
Kafka 最大的价值是削峰填谷和消费解耦。以日志场景为例:业务高峰时日志产生速度可能是平时的 5 倍,如果 Flume 直接写 HDFS,下游一旦反压,采集端就会堆积甚至丢弃;中间加一层 Kafka,生产端只管往 topic 里写,消费端按自己的速度慢慢拉,天然就把颠簸吃掉了。另外,数据要多方消费时(数仓一份、实时指标一份、安全审计一份),Kafka 的 fan-out 也最方便。
但 Kafka 不是越大越好。我见过一个项目,Kafka 集群 topic 建了上百个,每个 topic 还要 3 副本,新人一上来先被 topic 命名和权限管理搞晕。数据在 Kafka 里只应当作“暂存”,保留几天就够了,长期数据还是要落到 HDFS/数据湖。如果把 Kafka 当成数据仓库用,存储成本、topic 膨胀、分区不均,都会让你后期的维护量成倍上升。
举一个直接可抄的 Flume Kafka Sink 配置,就是把采集日志的挂载点接到 Kafka:
agent.sources = tail agent.channels = kafkaChannel agent.sinks = kafkaSink agent.sources.tail.type = spooldir agent.sources.tail.spoolDir = /data/logs agent.sources.tail.channel = kafkaChannel agent.channels.kafkaChannel.type = memory agent.channels.kafkaChannel.capacity = 10000 agent.sinks.kafkaSink.type = org.apache.flume.sink.kafka.KafkaSink agent.sinks.kafkaSink.kafka.bootstrap.servers = kafka-01:9092,kafka-02:9092 agent.sinks.kafkaSink.kafka.topic = raw-log agent.sinks.kafkaSink.channel = kafkaChannel这段配置只是示意,实际生产我会把 memory channel 换成 file channel 或 Kafka channel,防止 Flume 进程重启导致内存里的数据丢失。这个细节放到后面“翻车环节”展开。
3. 把需求翻译成参数:一张表完成采集方案的量化对比
工具认识得差不多,回到最核心的问题:怎么给团队一个可执行的结论,而不是继续在会议室里扯皮。我的做法是把需求量化成四个维度,给每个候选方案打分,并且把结论锁定在一个具体的业务场景上。“高可用”“高性能”这种词没有意义,“能扛 200MB/s 峰值、允许最多 1 分钟延迟、可容忍十万分之一丢弃率”才有讨论价值。
3.1 四个硬指标:吞吐量、端到端延迟、可靠性、人天成本
- 吞吐量。先说峰值,不要说均值。峰值决定了采集集群的规模。一条日志平均几百字节,10 万条/秒大约对应几十 MB/s;如果源头是点击流埋点,峰值可能是平均的 5 倍以上,这个系数必须打进冗余里。
- 端到端延迟。指从业务产生数据到数据可被查询的时间,不是某个组件的处理时间。要算上采集、消息排队、ETL、写入目标的时间。大多数场景做到分钟级,足够了。
- 可靠性。核心是你能不能接受丢数据。金融、订单、账户类不用想,肯定不能丢;埋点日志可以接受一定比例丢失,但这个口子一开,后面质量问题会很难受。建议用 at-least-once 加下游去重,而不是用 at-most-once 求快。
- 人天成本。包含开发调试和长期运维。一个只有 3 个数据工程师的团队,贸然引入 Flink CDC 整库同步,光是状态后端调优、DDL 变更处理就能让你怀疑人生。选型时一定要给“学习和排障成本”算一笔账。
这四个维度不是等比关系。场景不同,权重完全不一样。实时风控链路里延迟权重极高,离线条里可靠性优先级更靠前,小团队场景里人天成本往往一票否决。所以不要迷信“某某公司就是这么做的”——你并不清楚人家技术团队多厚,预算多高。
3.2 一张打分表:四个候选方案直接对比
假设你现在要采集电商平台的用户行为日志,日均日志量 20 亿条,峰值 5 万条/秒,业务方要求数据产生后 5 分钟内可查,团队有 4 个工程师。看四个候选方案的打分:
| 候选方案 | 吞吐量 | 延迟 | 可靠性 | 人天成本 |
|---|---|---|---|---|
| A. Filebeat + Kafka + Spark Streaming | 5 | 4 | 4 | 4 |
| B. Filebeat + Kafka + Logstash + Elasticsearch | 3 | 4 | 4 | 3 |
| C. Flume + HDFS(离线) | 3 | 2 | 3 | 4 |
| D. Flink CDC + Kafka + Flink | 4 | 5 | 5 | 2 |
按 1 到 5 打分,5 表示最满足。这张表的价值不在于数字精确,而在于把争论从“我觉得 A 好”变成“延迟这项你打 3 的理由是什么”。我们团队实际选的是方案 A,理由很简单:延迟 5 分钟内完全达标,吞吐和可靠性足够,团队在 Spark 上已有基础,不需要为了实时而扩大技术面。
你也可以把这四个维度设计成加权评分:先把业务方对延迟、可靠性的底线写死,在满足底线的方案里再挑吞吐和人天成本最优的。重点是“先排除,再优选”,而不是从零开始比好坏。
3.3 三种典型业务场景的推荐组合
最后套三个具体场景,方便你对号入座。
场景一:电商 / 在线教育埋点日志。数据源是 Web 和 App 埋点,量级大、字段乱、实时性要求中等。推荐组合是 Filebeat 采集,进 Kafka 按业务分 topic,再由 Flink/Spark 清洗后落 HDFS/ES。关键点是埋点原始数据在 Kafka 原始 topic 里保留一份,清洗后的数据再进数仓,方便日后重算。
场景二:金融 / 交易系统库表同步。源库是核心业务的 MySQL,要求分钟级同步到数仓,还要保留变更历史。推荐组合是首次全量用 DataX 迁底,增量用 Flink CDC 监听 binlog 投到 Kafka,再由下游 Flink 做 join 和写入。全量和增量之间要处理好水位线衔接,否则会出现“全量已经跑完但 binlog 从更早位置开始”的重叠或空洞,这个话题第四部分还会提。
场景三:车联网 / 工业物联网设备数据。设备产生的是高频小幅消息,网络抖动多,无法直接跑在公网。推荐组合是设备经 MQTT 网关接入(EMQX/Mosquitto),再由网关侧 Kafka Connect 把数据转存到 Kafka,之后按实时/离线分别交给 Flink 和 HDFS。这里不要试图让采集 agent 直接连设备,协议兼容和弱网处理不是你该在数据层解决的问题。
4. 方案落地后最容易翻车的四个环节:别让选型毁在上线后
方案选定、任务排期、ETL 都没问题——这时候最容易松劲。实际上,我看到的大多数采集链路问题,不是选型选错,而是上线后的运维细节没跟上。下面四个坑,基本可以覆盖 80% 的“采集链路事故现场”。
4.1 丢数据:先搞清楚是哪个环节吞了你的数据
丢数据的常见位置有三处:采集 agent 崩溃、消息队列限流、目标端写失败。
- agent 端:Flume 如果用 memory channel,JVM 一重启,缓冲在内存里没来得及投递的数据全没了。稳妥做法是换 file channel 或 Kafka channel,用磁盘换内存,牺牲一点吞吐换不丢。Filebeat 相对安全,它用 registry 文件记录偏移量,但也要注意磁盘坏块和日志轮转配置的配合。
- 消息队列端:Kafka 的
acks=0或acks=1在极端情况下会丢消息,生产环境至少要acks=all。另外,broker 端unclean.leader.election.enable如果设成 true,主副本缺失时选出来的节点可能没有最新数据。这个参数在很多默认配置里挺坑的,不提前设好,故障切换时就会悄悄丢数据。 - 目标端:HDFS 写失败最常见的不是网络,而是小文件过多引发 name node 压力,或者分区目录权限问题。采集任务最好带自动重试,并对连续失败次数做告警,不能重试了就默默跳过。
出问题后不要拍脑袋“数据量不大,丢了就丢了”。先核对三处:source 端有没有记录已读位置、Kafka 端每条消息的 offset 范围、sink 端实际写入条数。三角对不上,说明链路某个环节的语义并不是你以为的 at-least-once,这时候要回到配置层面逐段排查。
4.2 重复数据与幂等消费:至少一次和恰好一次的区别
选择 at-least-once,意味着重复几乎不可避免。Flume 重启后可能重新读取最后一批日志,Filebeat 重传了没确认完成的一段,下游统计时如果不做去重,报表就会偏高。去重思路要提前设计,别等出问题再补。
- 给每条数据一个稳定的事件 ID(例如 UUID 或者日志自带 requestId),下游写入目标表时用事件 ID 做唯一键;
- Flink 里可以用 keyby 事件 ID + 状态去重,或者依赖目标存储的幂等写入;
- 如果数据要落到 HDFS,按业务时间去重分区,下游读时对重复部分取最早一条即可。
最容易忽略的是“全量 + 增量”衔接造成的重复或漏数。DataX 先做全量,同时 binlog 已经开始捕获增量,如果全量快照的时间点没有对齐,增量里很可能会包含全量已经搬过的数据。建议在全量开始时记录 binlog 位点,全量结束后,从该位点消费增量,并让下游按主键做 upsert。这样即使两边数据有重叠,upsert 也能把结果收敛到正确的最终形态。
4.3 Schema 变更:采集层必须提前设计应对策略
业务表加个字段、日志打印格式换了个分隔符,在采集层看来都是“事故”。如果不提前设计,任何上游改动都会传导到下游任务失败。这块我吃过两次亏,现在总结成三条经验:
- 原始层尽量存宽松格式。日志和 CDC 变更在原始层用 JSON 或 Avro 保留,不要一上来就强行转成定长结构。这样上游加字段,下游只要不强制解析,链路就不会断。
- 引入 Schema Registry(例如 Confluent Schema Registry)管理 Avro/JSON Schema,版本升级时做兼容性校验,加字段要带 default 值。
- 建立 DDL/DML 变更的报备流程。业务方改表结构之前,先同步给数据团队,采集任务、目标表、下游模型一起评估。
我见过最痛的一个案例:业务在订单表上加了两个字段,因为上游工具没有版本管理,Flink CDC 任务直接挂掉,数据从凌晨断到中午,恢复之后还要补数。那次之后,我们把所有对接表的 schema 纳入版本管理,并加了字段变更的告警钩子,再没出过同类问题。
4.4 链路不能裸奔:采集层监控和团队能力底线
采集链路没有监控,等同于蒙眼开车。很多团队把监控资源都投在计算引擎和数仓,忽略了“数据进没进来”这个最上游的问题。我建议至少盯住四个指标:
- 采集 agent 的发送速率和 source 读入速率,如果两者开始拉开,说明 agent 在积压;
- Kafka 消费组的 lag,lag 持续上涨,说明消费端跟不上;
- sink 的写失败次数和重试次数,任何一个不为 0 都该有告警;
- 端到端的数据量基线,例如每 10 分钟自动对比今日峰值和昨日同期,偏差超过 30% 就告警,这是抓“上游静默不产数据”最有效的办法。
再说团队能力底线。一个技术方案再完美,如果团队里没人愿意长期维护,上线三个月后它就会变成新的技术债。选型时请诚实地评估:你们能不能处理 Flink checkpoint 失败?能不能调 Kafka 分区和参数?遇到 binlog 解析位点偏移,有没有人能快速定位?如果答案不乐观,宁可先用简单方案把业务跑起来,再逐步演进到更复杂的链路。技术在不停迭代,但团队对复杂度的消化能力是有上限的,选型要为人服务,不是为简历服务。