5.4 数据进出集群:Sqoop与Flume


文档摘要

5.4 数据进出集群:Sqoop 与 Flume 本节摘要:Sqoop 用 MapReduce 并行搬运关系库与 HDFS 之间的批量数据,Flume 用"源-通道-汇"三段式结构持续收集日志流进 HDFS。一个管批量管道、一个管流式水龙头,两者构成数据进出 Hadoop 集群的南北通道。本节给出两者的配置实战与常见故障定位。 一进一出,两种节奏 数据旅程的入口有两类节奏。批量节奏:业务库每天凌晨导出昨日订单,整表或增量切片进 HDFS——量大、定时、可重跑,Sqoop 的领地。流式节奏:应用服务器源源不断产生日志,希望尽快可查——量大、持续、不可重放(丢了就是丢了),Flume 的领地。

5.4 数据进出集群:Sqoop 与 Flume

本节摘要:Sqoop 用 MapReduce 并行搬运关系库与 HDFS 之间的批量数据,Flume 用"源-通道-汇"三段式结构持续收集日志流进 HDFS。一个管批量管道、一个管流式水龙头,两者构成数据进出 Hadoop 集群的南北通道。本节给出两者的配置实战与常见故障定位。

一进一出,两种节奏

数据旅程的入口有两类节奏。批量节奏:业务库每天凌晨导出昨日订单,整表或增量切片进 HDFS——量大、定时、可重跑,Sqoop 的领地。流式节奏:应用服务器源源不断产生日志,希望尽快可查——量大、持续、不可重放(丢了就是丢了),Flume 的领地。出口侧类似:分析结果批量回写关系库供报表使用(Sqoop 导出),或实时入 HBase 供在线服务(Flume 汇端可写 HBase)。

判断用哪个工具的问题等价于判断数据的节奏:能攒批就 Sqoop,必须边产生边搬运就 Flume。混用是最常见的架构错误——用 Flume 搬每日快照(重跑困难、无事务边界),或用 Sqoop 搬秒级日志(数据库根本扛不住这种轮询)。

Sqoop:把 JDBC 搬运并行化

Sqoop(SQL-to-Hadoop)的本质:读源库的表元数据,生成一个只有 Map 的分布式作业,按主键区间把表切成 N 段,N 个 Map 各自通过 JDBC 拉一段、直接写 HDFS。

sqoop import \ --connect jdbc:mysql://db.example.com:3306/orders \ --username etl --password-file /user/etl/.dbpass \ --table t_order \ --columns "order_id,user_id,amount,created_at" \ --where "created_at >= '2026-08-18' AND created_at < '2026-08-19'" \ --target-dir /warehouse/ods/t_order/dt=2026-08-18 \ --split-by order_id -m 8 \ --fields-terminated-by '\t' # -m 8:8个并行Map --split-by:切分列(无主键表必须显式指定)

几条实战要点。切分列必须有均匀分布:按主键或递增列切 8 段,若某段占 90% 行(倾斜),并行等于白搭。增量导入用 increment 模式(append 按自增列 / lastmodified 按时间戳列),配合 --last-value 由调度层记住断点。密码永远用 password-file,命令行明文密码会进 shell 历史与 YARN 日志。导出方向(HDFS → 关系库)同样是 Map 作业并行 insert/update,注意目标表先清空或用 staging 表加事务切换,防止半程失败留下脏数据。

数据类型映射是隐性深坑:源库 DECIMAL(38,10) 映射不足会截精度;时间字段时区差异让 dt 分区错一天。Sqoop 提供 --map-column-java 手工指定映射,重要管道上线前必须用小样本对账(行数、金额总和两个数即可拦住绝大多数映射事故)。

Flume:源、通道、汇

Flume 的模型只有三个角色:Source(收数据:盯日志文件目录、开端口收事件、接上游 Agent)、Channel(暂存:内存通道快但断电丢、文件通道慢但可靠)、Sink(发出:写 HDFS、转发下一跳、写 HBase/Kafka)。Agent 是三者的组合进程,多级 Agent 串联即成采集链。

应用服务器 集群边缘 HDFS tail日志 → Agent1 ──Avro──▶ Agent2 ──▶ HDFS Sink(按时间/大小滚动分区文件)
# Agent2 核心配置 边缘汇聚节点 a2.sources = r1 a2.channels = f1 a2.sinks = h1 a2.sources.r1.type = avro a2.sources.r1.bind = 0.0.0.0 a2.sources.r1.port = 4141 # 文件通道:落盘暂存 上游崩了数据也不丢 a2.channels.f1.type = file a2.channels.f1.checkpointDir = /data/flume/checkpoint a2.channels.f1.dataDirs = /data/flume/data a2.sinks.h1.type = hdfs a2.sinks.h1.hdfs.path = /data/access/dt=%Y%m%d/hour=%H a2.sinks.h1.hdfs.filePrefix = access a2.sinks.h1.hdfs.rollInterval = 3600 a2.sinks.h1.hdfs.rollSize = 134217728 a2.sinks.h1.hdfs.rollCount = 0 a2.sinks.h1.hdfs.fileType = DataStream a2.sinks.h1.hdfs.round = true a2.sinks.h1.hdfs.roundValue = 5 a2.sinks.h1.hdfs.roundUnit = minute

读这份配置要抓四个设计决策。通道选 file 不选 memory:日志不可重放,可靠性优先,文件通道的事务机制保证 Source 收到 Sink 发出之间不丢。HDFS 滚动三条件(时间 1 小时、大小 128MB、条数禁用):同时命中任一即切文件——目标是让文件接近块大小,正对 2.4 节小文件治理。round 取整到 5 分钟:时间戳取整让路径稳定,避免事件乱序造成的目录碎片。fileType 用 DataStream:文本直写,若下游要 ORC 应在汇聚后由离线任务转换,而不是在 Sink 里做重格式化。

图 5-4 数据进出集群的两条通道

图 5-4 数据进出集群的两条通道

故障定位:两个高频事故

Sqoop 空跑或慢跑:作业成功但目录为空——多半 --where 条件写错字段或时区差一天;作业成功但只有 1 个 Map 实际干活——切分列倾斜,看各 Map 的记录数计数器即可确认。

Flume 小文件淹没:HDFS Sink 目录里出现海量几 KB 文件,直接病根是滚动三条件全没命中正确值——rollInterval=0 且 rollSize 过小,或上游流量骤降让文件长期敞开不切。症状在 2.4 节的账本里已经写过:NameNode 内存告警、Map 作业分片爆炸。修法按"目标 128MB"重算三条件,必要时加 hdfs.batchSize 控制冲刷频率。另一个反向事故是文件长期不关闭(路径里 hour=HH 的文件一直 OPEN),通常是 Agent 时钟漂移或 round 配置错误,下游 Hive 查不到最新分区的第一嫌疑人就是它。

本节要点回顾

  • 节奏决定工具:能攒批的走 Sqoop 并行 JDBC 管道,持续日志走 Flume 源通道汇链路;
  • Sqoop 是纯 Map 作业:按切分列分段并行,切分列倾斜则并行失效;增量靠断点模式;导出用 staging 表保事务边界;
  • 类型映射要对账:精度与时区是两大隐性坑,行数加金额总和的小样本对账拦住九成事故;
  • Flume 通道选 file:暂存落盘换来不丢数据;HDFS 滚动条件对齐块大小,防小文件淹没;
  • 两条铁律:Sqoop 慢先看分段计数器,Flume 乱先看文件滚动状态。

下一节把前面所有站点连成自动流水线:Oozie 工作流。


作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U