ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Sqoop+Oozie实战:构建MySQL到Hive的自动化数据迁移调度系统

Sqoop+Oozie实战:构建MySQL到Hive的自动化数据迁移调度系统 这阵子刚好有个数据仓库的活每天晚上要从十来个MySQL业务库往Hive里同步数据。刚开始图省事写了个shell脚本挂crontab直接跑sqoop命令结果跑了不到一个月就翻车了——不是某张表新增了字段导致导入报错就是上游凌晨有个批量任务把源库拖慢了整个同步链路全卡在一起大早上看着告警群里一堆失败任务头皮发麻。后来把Sqoop和Oozie拎出来重新捋了一遍搭了一套自动化数据迁移调度系统才算是把这块彻底收拾利索。如果你也在折腾这类场景——用Sqoop做关系库和Hadoop之间的数据迁移用Oozie做定时调度和工作流编排那这篇实战记录应该对你有用。我会从方案选型一直讲到我实际踩过的各种坑配置、命令、踩坑过程都放出来尽量让照着做的人能少走弯路。1. 为什么是Sqoop Oozie这套组合到底解决了什么1.1 数据迁移的痛点不止是“把数据搬过去”先说个容易被低估的问题数据迁移看起来就是把一张表从MySQL挪到Hive跑一条sqoop import就完事了。但放到实际业务里麻烦事全在被忽略的部分。比如一张订单表每天新增几百万行是全表覆盖还是增量追加源库表结构半夜被DBA调整了目标表怎么自动感知迁移任务挂在凌晨两点失败了是重跑还是跳过这些要是都靠人肉盯那这份工作基本就没法干了。再往深一层说数据迁移是有依赖关系的。事实表要在维度表刷完之后才能做关联清洗ODS层的表要先同步完DWD层的加工任务才能启动。这种跨任务、跨时间、带重试和告警的流程纯靠crontab加shell脚本几乎管理不过来——脚本里的判断逻辑越堆越多到最后你自己都分不清哪个任务失败会影响哪条链路。所以真正需要的不是一个导数据的工具而是一套能编排、能调度、能恢复的完整体系。Sqoop负责解决“怎么把数据搬过去”Oozie负责解决“什么时候搬、搬完干什么、失败了怎么办”俩搭一起就是一套能睡得着觉的自动化方案。1.2 Sqoop和Oozie各自擅长什么Sqoop这个东西用一句话说就是关系型数据库和Hadoop之间的搬运工。它本身分成Sqoop 1和Sqoop 2两个版本但到2024年了Sqoop 2还是那个半死不活的状态生产环境里绝大多数人用的都是Sqoop 1这套架构简单直接——命令行提交作业底层把导入导出任务翻译成MapReduce执行。关于版本我先劝一句别碰Sqoop 2社区基本不维护坑远比1多。Sqoop之所以能在数据迁移这块站稳脚跟核心在于它理解关系库的那套东西。比如导入一张MySQL表它会先通过JDBC读取表结构把字段映射成Hive表的列然后根据主键或者指定的split列去计算数据分布生成多个map任务并行拉取数据。这不光是快更重要的是它对数据库方言的适配做得不错——MySQL、Oracle、PostgreSQL、SQL Server都有对应的连接器字段类型转换、主键识别、blob/clob处理这些脏活都被它包了。Oozie则是Hadoop平台上的工作流调度引擎。它的玩法是把任务定义成DAG有向无环图一个节点跑完知道下一个节点该干啥节点之间还能传数据、传参数。它不支持环——官方就不允许循环依赖这样反而逼着你把业务流程理顺了再落地。Oozie里干活的基本单位是action一个action可以是一个Sqoop命令、一个Hive脚本、一个Shell脚本、甚至一个Java程序action之间用transition串起来。再加上Coordinator的定时触发机制这套东西就能做到“每天凌晨2点自动跑全量同步 → 同步成功才触发数据清洗 → 清洗完给下游发通知”。1.3 这套方案和别的工具对比有什么优势我知道肯定有人想说现在DataX、Canal、NiFi都挺火的为什么非得用Sqoop加Oozie。说实话这取决于你所在的平台环境。如果你的Hadoop集群是CDH或者HDP发行版那Sqoop和Oozie全是自带的组件不需要额外部署东西版本兼容性也是厂商测过的。用DataX还得自己弄个调度平台去管它。选择什么方案首先要看你的底座提供了什么而不是哪个工具听起来更时髦。另外Oozie最值钱的一个特性是它和Hadoop生态深度绑定——提交作业本质上是往YARN上丢Application资源管理、日志聚合、容错重试都是现成的。白天集群跑实时任务晚上跑批处理Oozie作业能自动排队等资源不会把集群打爆。这点用Quartz或者XXL-Job这类外挂调度工具是做不到的它们只管触发命令管不了跑起来的MapReduce到底吃多少资源、卡住了怎么处理。一句话总结我的选型逻辑如果集群本身就带着Sqoop和Oozie别犹豫直接用如果是从零搭这两种依然是稳定性和社区支持最均衡的组合。2. 环境准备与Sqoop基础实操2.1 搭建这套东西需要什么底子先说下我这边的基础环境方便你对号入座。用的是CDH 6.3.x发行版Hadoop 3.0.0Hive 2.1.1Oozie 5.1.0Sqoop 1.4.7。TaoBao出的那个Sqoop也有人在用但我在生产环境没试过就不胡乱推荐了。如果是从Apache原版自己搭有几样东西是绕不开的Hadoop集群至少NameNode加DataNode能正常跑MapReduce、YARN资源调度正常、Hive Metastore服务在线、Oozie服务端部署好并且能连上YARN和HDFS。Sqoop装在能提交作业的节点上一般是集群的客户端节点或者边缘节点。还有最关键的一样东西——MySQL的JDBC驱动。Sqoop本身不带MySQL驱动得自己下载mysql-connector-java jar包丢到Sqoop的lib目录里。这个坑我见过好几个人踩过sqoop命令报了错日志里明晃晃写着java.lang.ClassNotFoundException: com.mysql.cj.jdbc.Driver结果折腾半天就是驱动没放对位置。版本上有个小细节驱动一定要和MySQL服务端版本匹配不用完全一致但主版本不能差太远。我这边MySQL是5.7用的驱动是5.1.49跑了快两年没出过兼容性问题。如果你的MySQL是8.0最好直接用8.0.x的驱动否则可能会碰上一堆认证插件和时区的问题。2.2 连接MySQL的关键配置从URL到连接参数Sqoop导入MySQL数据的第一步就是写对JDBC连接串。这个连接串看着简单但有几个参数在实际生产中必须带上否则早晚会有坑jdbc:mysql://192.168.10.20:3306/business?useSSLfalsecharacterEncodingutf8serverTimezoneAsia/ShanghaiuseSSLfalse是必须的——开发环境或者内网环境MySQL如果没开SSL加上这个参数反而是好事不然连接握手阶段会莫名卡住报SSL握手失败。这里再多说一句如果是内网大数据集群数据链路本身在可控网络内关掉SSL不仅为了省事还能减少握手时的资源消耗。characterEncodingutf8必须显式指定。不然你在sqoop命令行里写的查询条件如果带中文大概率会出现乱码或者查不到数据的情况。serverTimezoneAsia/Shanghai是给MySQL 8.0驱动用的。MySQL 8.0以上版本的驱动对时区很敏感不指定的话驱动会去读系统默认时区如果数据库服务器和你的应用服务器时区不一致timestamp类型的字段导出来会差8个小时。修数据可比写代码痛苦多了。连接字符串还有一个容易被忽视的点JDBC URL里的主机名到底能不能被解析。我见过有人在连接串里写了内网IP本机ping得通但Sqoop提交的MapReduce任务跑在数据节点上那些机器未必能访问这个IP。正确做法是在每台DataNode的/etc/hosts里都配上源库的映射或者统一用内网域名别赌IP在全网通。这个问题排查起来非常恶心——因为在本机测试一切正常一上集群就报连接超时。Sqoop的命令行参数里还有几个和连接相关的关键项顺带记一下--connectJDBC连接串上面详述--username/--password数据库账号密码密码是明文参数生产环境建议用--password-file从HDFS读否则ps能看到密码--driver手动指定驱动类名Sqoop有时候会根据URL自动推断但显式写出来更稳妥2.3 从零跑通第一条导入命令先把最简单的全量导入跑通验证环境和权限都没问题。假设我们有一张MySQL的business.orders表要导入到Hive的default.orders表sqoop import \ --connect jdbc:mysql://192.168.10.20:3306/business?useSSLfalsecharacterEncodingutf8serverTimezoneAsia/Shanghai \ --username bigdata \ --password-file /user/bigdata/mysql.pwd \ --table orders \ --hive-import \ --hive-table default.orders \ --create-hive-table \ --m 4这条命令里值得注意的几件事--m 4是map并发数这个直接决定导入速度。4个并发不算激进但也足够跑通大部分场景。map数不是越大越好得看源库的负载和这张表的数据量后面优化部分我会细说。--hive-import会先把数据从MySQL拉到HDFS临时目录再执行Hive的load data命令把临时目录的数据load进目标表。中间多了一步但好处是Hive表的元数据是自动生成的不需要你手动去写CREATE TABLE。--create-hive-table的意思是目标Hive表如果不存在就自动创建。这个参数在生产环境要慎用——如果表已经存在加上这个参数会直接drop掉重建后果自己体会。我第一次跑测试的时候就是因为表已经建好了又习惯性带了这个参数结果整个表被清空重建幸好是测试环境。跑完去Hive里看一眼数据量对不对再查一下HDFS上实际生成的文件数和大小确认4个map都分片成功。2.4 增量导入和参数调优别再用全量覆盖了全量导入通了之后马上要考虑增量场景。正常业务表动辄几百万上千万行天天全量覆盖既不经济也容易出问题。Sqoop提供两种增量导入模式--incremental append适合表里有自增主键的场景每次只导主键值比上次最大值大的数据。第一次跑的时候用--last-value 0之后每次跑之前要查一下上次的last-value存哪了——这有点绕正常做法是把last-value记录到一张管理表里或者直接用Oozie的coordinator参数动态传进去。--incremental lastmodified适合表里有一个最后修改时间字段的场景。这种模式会按时间字段过滤出最近改过的数据比append更精确但要求表里必须有维护良好的时间列而且如果数据是在导入过程中被修改的可能会漏数据。增量导入有一个隐藏的大坑并行度。如果增量数据量很大比如一次导入上千万行map数必须跟着涨上去。但map数涨了之后如果split列增量时一般用主键做split数据分布不均匀个别map还是会跑得很慢。解决方案是配合--boundary-query手动指定split边界让Sqoop用你自己定义的查询去分段而不是傻乎乎地均匀切分。比如订单表的主键是order_id但分布严重不均——早期订单稀疏近期订单密集这时候可以写--boundary-query SELECT MIN(order_id), MAX(order_id) FROM orders WHERE update_time 2024-06-01让spilt的范围收敛到增量数据实际所在的主键区间避免把一整年的主键区间全扫描一遍。另外一个影响性能的重要参数是--fetch-size它让每个map批次从数据库下拉的行数。默认值可能不太适合高延迟的数据库链路调成1000或者2000通常能明显提升大表导入的吞吐量。注意别调太大否则内存会爆JVM直接OOM。3. Oozie调度工作流的设计与搭建3.1 Oozie工作流的基本结构先用一段大白话解释Oozie的工作流到底长啥样。它就是一堆XML文件核心是一个workflow.xml描述了整个任务链?xml version1.0 encodingUTF-8? workflow-app xmlnsuri:oozie:workflow:0.5 namemysql-to-hive start tosqoop-orders/ action namesqoop-orders sqoop xmlnsuri:oozie:sqoop-action:0.3 job-tracker${jobTracker}/job-tracker name-node${nameNode}/name-node configuration property namemapred.job.queue.name/name value${queueName}/value /property /configuration commandimport --connect ${jdbcUrl} --username ${jdbcUser} --table orders .../command /sqoop ok toend/ error tosend-alert/ /action action namesend-alert shell xmlnsuri:oozie:shell-action:0.2 job-tracker${jobTracker}/job-tracker name-node${nameNode}/name-node execbash/exec argumentsend_alert.sh/argument /shell ok toend/ error toend/ /action end nameend/ /workflow-app几个关键点start to.../和end name.../是所有工作流必须的起止节点节点名称自己定义但别重复。action里面放实际要跑的任务类型Sqoop action有专门的名字空间uri:oozie:sqoop-action:0.3版本号和Oozie版本有关5.x基本都用0.3。ok to.../和error to.../是两个固定分支任务成功走ok失败走error。这是整个工作流里处理异常全靠的地方——通常error分支接一个发告警的shell动作或者接一个重试逻辑。之前见过新手在Oozie里写Sqoop命令时踩过一个坑把command换成了arg然后一个参数一个参数地单独指定结果发现总有几个参数传不进去报奇怪的语法错误。这里明确说一下command是Oozie的Sqoop action指定的标准写法它的内部实现会把整个字符串按空格切成参数数组像普通的Sqoop命令一样解析。如果你在命令行里加了引号比如查询条件--where id 100到Oozie里写在这个标签里的引号就得去掉否则会被当成参数的一部分导致SQL语法报错——这是个极其经典的坑我在这上面浪费过一下午。3.2 Coordinator定时调度的核心配置光有workflow还不够它只是定义“怎么干”没定义“什么时候干”。定时调度的活是Coordinator干的配置在coordinator.xml里?xml version1.0 encodingUTF-8? coordinator-app xmlnsuri:oozie:coordinator:0.4 namemysql-sync-coord frequency${coordFrequency} start${startTime} end${endTime} timezone${coordTimezone} xmlns:slauri:oozie:sla:0.2 action workflow app-path${workflowAppPath}/app-path configuration property namejobTracker/name value${jobTracker}/value /property property namenameNode/name value${nameNode}/value /property property namequeueName/name value${queueName}/value /property property namejdbcUrl/name value${jdbcUrl}/value /property property namejdbcUser/name value${jdbcUser}/value /property property namelastValue/name value${lastValue}/value /property /configuration /workflow /action /coordinator-app这个配置里最重要的是frequency和start/end。frequency表示运行频率单位是分钟比如frequency1440就是每天跑一次。start和end是时间窗口规定这个调度任务在哪个时间段内生效格式类似2024-01-01T00:00Z这里时区一定要想清楚我一般用timezoneAsia/Shanghai显式指定否则默认是UTC你定了凌晨2点跑实际可能是北京时间早上8点才跑白等半宿。Coordinator还有个很爽的特性它可以在每次触发时生成一组变量比如${coordinator:actualTime()}可以拿到真正的执行时刻。这个经常拿来做增量迁移的时间水位——比如你要每天同步前一天的数据可以直接在workflow里拼SQL--where create_time ${coordinator:dateOffset(coordinator:actualTime(), -1, DAY)}这样连外部传参都省了调度系统自己算出来该同步哪天的数据。还有一个组件叫Bundle它是Coordinator的集合可以把多个Coordinator打包在一个Bundle里统一启停。如果你的同步任务涉及多个表、多套coordinator用Bundle管理会方便不少。3.3 工作流的部署与提交Oozie的工作流不是一个本地文件它需要把整个目录上传到HDFS然后通过Oozie服务端去调度。目录里至少要放workflow.xml、coordinator.xml和一个job.properties文件后者的作用是定义所有用到的变量的值nameNodehdfs://nameservice1 jobTrackerresourcemanager:8032 queueNamedefault workflowAppPath/user/bigdata/oozie/worksflows/mysql-sync jdbcUrljdbc:mysql://192.168.10.20:3306/business?useSSLfalsecharacterEncodingutf8serverTimezoneAsia/Shanghai jdbcUserbigdata oozie.wf.application.path${nameNode}/user/bigdata/oozie/worksflows/mysql-sync oozie.coord.application.path${nameNode}/user/bigdata/oozie/worksflows/mysql-sync然后按这个顺序操作# 先把整个工作流目录推到HDFS hdfs dfs -rm -r /user/bigdata/oozie/worksflows/mysql-sync hdfs dfs -put ~/oozie-wf/mysql-sync /user/bigdata/oozie/worksflows/ # 提交coordinator任务 oozie job -oozie http://oozie-server:11000/oozie \ -config ~/oozie-wf/mysql-sync/job.properties \ -submit # 如果想提交后立即启动用 -run 参数 oozie job -oozie http://oozie-server:11000/oozie \ -config ~/oozie-wf/mysql-sync/job.properties \ -run-submit只是把作业登记到Oozie不会立即触发-run提交后立刻启动。第一次测试建议先用-run跑一次确认整个链路能通。之后再用-submit配合Coordinator的定时调度机制。部署里最容易踩的坑是HDFS的路径权限。Oozie服务端用户一般是oozie它会去读你上传的配置文件。如果目录权限设成700Oozie进程就可能没权限读取报File does not exist或者其他千奇百怪的错误。建议统一用hdfs dfs -chmod -R 755把目录放开或者把目录owner改成oozie用户。另一个常态问题每次改了workflow.xml都要重新put一遍到HDFS然后让Oozie重新加载。就算你只改了一个参数Oozie也不会自动感知你本地文件变了必须重新上传。我则习惯写成一个小脚本一键完成rm、put、重新提交三步避免手动操作出错。3.4 Oozie与Sqoop的衔接细节在Oozie里跑Sqoop action最重要的就是搞明白变量传递的机制。job.properties里的键值对会通过${}占位符替换到workflow.xml和coordinator.xml里替换发生在Oozie服务端不是发生在客户端。这意味着你在workflow.xml里写的${jdbcUrl}、${lastValue}最终都会由Oozie运行时去job.properties对应的配置里取值。有个容易忽略的细节Oozie的Sqoop action里连接串如果包含特殊字符比如XML里必须转义成amp;。这个是最经典的踩坑点了——你在本地sqoop命令里写useSSLfalsecharacterEncodingutf8没问题但写到XML里不转义Oozie直接给你报XML解析错误整个作业还没开始跑就挂了。还有个问题是密码。在命令行里跑Sqoop可以直接写--password但在Oozie的workflow.xml里明文写密码一旦HDFS被打包下载或者同步到别的地方密码就泄露了。更稳妥的做法是让Sqoop任务从HDFS的密码文件读取--password-file hdfs://nameservice1/user/bigdata/oozie/mysql.pwd然后确保这个密码文件只有任务运行用户和oozie用户能读。注意--password-file期望的是HDFS上的路径和命令行里用本地文件系统的--password-file语义不太一样——在Oozie环境里一定要走HDFS路径否则每个map任务都会找不到文件报错。4. 从单表同步到全库调度一个能直接用的配置示例4.1 增量数据迁移的两种常用策略前面讲了Sqoop的append和lastmodified两种增量模式但要真正落地到自动调度里还必须想清楚两个问题增量标识从哪里来以及目标表怎么合并增量数据。先说增量标识。自增主键模式最简单但有个缺陷如果业务库发生主键回退或者人工删数据再插数据主键大小和实际修改时间就不完全对应了会漏数。lastmodified模式更准确但要求上游表必须有类似update_time或者modify_time的字段而且这个字段要真的随每次update更新——有些开发偷懒字段建了但没写代码维护结果每天增量都是全量。我吃过的亏提醒你接一张新表的时候先手动查一下UPDATE_TIME最多的那几条确认它是真的在变。再说目标表合并。千万不要以为Sqoop把增量数据load进Hive表就完事了那一堆增量文件会和历史全量文件混在一起。最常见的做法是三层结构ODS层用分区表每天一个分区增量数据导入到dt2024-06-15这个分区。这本身没毛病。但DWD层如果要的是最新全量就得每天做一次INSERT OVERWRITE TABLE ... SELECT ... FROM ODS WHERE dt today把历史数据攒一遍。用Oozie编排时一个Sqoop action导入ODS紧接着一个Hive action跑合并SQL数据链路清晰重试成本也低。下面是一套按日分区的全流程工作流先判断增量是否为空不为空再跑合并避免空数据把全量覆盖掉。这一步看着简单其实相当关键我进过的项目里出现过好几次空增量覆盖全量的事故——因为上游archive策略不同某个分区的数据源恰好没生成文件sqoop导入成功但分区空而后面的合并SQL没判断就执行了最后全表被清空。4.2 一个生产可用的完整工作流配置包目录结构长这样/user/bigdata/oozie/worksflows/mysql-sync/ ├── workflow.xml ├── coordinator.xml ├── job.properties └── lib/ └── mysql-connector-java-5.1.49.jar注意lib目录这是另一个坑——Sqoop action运行的时候JVM要用到MySQL驱动但这个驱动不一定会出现在所有节点的classpath里。Oozie的workflow目录下放一个lib并把驱动jar丢进去Oozie会把lib目录里的jar自动加到classpath里省掉在每个节点上装驱动。这个技巧能避免很多莫名其妙的ClassNotFoundException。workflow.xml核心配置完整版拿过来改改就能用?xml version1.0 encodingUTF-8? workflow-app xmlnsuri:oozie:workflow:0.5 namemysql-orders-sync start tosqoop-orders/ action namesqoop-orders sqoop xmlnsuri:oozie:sqoop-action:0.3 job-tracker${jobTracker}/job-tracker name-node${nameNode}/name-node configuration property namemapred.job.queue.name/name value${queueName}/value /property property namemapreduce.map.memory.mb/name value2048/value /property /configuration commandimport --connect ${jdbcUrl} --username ${jdbcUser} --password-file ${passwordFile} --table orders --target-dir /user/hive/warehouse/ods.db/orders/dt${syncDate} --incremental lastmodified --check-column update_time --last-value ${lastValue} --split-by order_id --m 6 --fetch-size 1000 --fields-terminated-by \001 --null-string \\N --null-non-string \\N/command /sqoop ok tohive-merge/ error tocheck-error/ /action action namecheck-error shell xmlnsuri:oozie:shell-action:0.2 job-tracker${jobTracker}/job-tracker name-node${nameNode}/name-node execbash/exec argument/tmp/alert.sh/argument argumentmysql-sync-workflow/argument argumentsqoop-orders/argument /shell ok toend/ error toend/ /action action namehive-merge hive xmlnsuri:oozie:hive-action:0.5 job-tracker${jobTracker}/job-tracker name-node${nameNode}/name-node configuration property namemapred.job.queue.name/name value${queueName}/value /property /configuration scriptmerge_orders.sql/script paramdt${syncDate}/param /hive ok toend/ error tocheck-error/ /action end nameend/ /workflow-app几个容易出问题的点专门提醒一下--target-dir这里没有走--hive-import而是直接把数据写到HDFS上的Hive表目录里。这样做的好处是省掉SQOOP自动load数据到Hive表的那一步后面用Hive SQL自己管理分区更灵活。但代价是你必须保证目标目录的路径和Hive表的分区结构完全对得上名字怎么拼都不能错。我建议你自己搭好一套约定比如所有ODS表的目录都用/user/hive/warehouse/ods.db/{表名}/dt{日期}这样就算换表也不会迷路。--fields-terminated-by \001和--null-string \\N这四个参数是配套的作用是把数据字段和空值都固定成Hive默认的格式。不然Hive建表默认的分隔符是\001空值默认是NULL而Sqoop默认的字段分隔符可能是逗号空值可能是个空字符串——你导入之后Hive表里看起来数据是有的但一切换格式查询就全乱套。--null-string是给字符串类型的空值用的--null-non-string是给数值类型用的。两个都设成\\N注意这里是两个反斜杠是为了让Sqoop把数据库里的NULL转成Hive能识别的\N文本表示否则Hive表查出来都是空但底层文件里其实写了一堆null字符串——这个坑会直接影响下游的聚合结果。4.3 Coordinator的定时配置和参数传递coordinator.xml里的关键设置coordinator-app xmlnsuri:oozie:coordinator:0.4 nameorders-sync-coord frequency${coordFrequency} start${coordStart} end${coordEnd} timezoneAsia/Shanghai action workflow app-path${workflowAppPath}/app-path configuration property namesyncDate/name value${coordinator:dateOffset(coordinator:actualTime(), 0, DAY)}/value /property property namelastValue/name value${coordLastValue}/value /property ... /configuration /workflow /action /coordinator-app这个${coordinator:dateOffset(coordinator:actualTime(), 0, DAY)}是Oozie的EL表达式作用是把当天的日期作为参数syncDate传给workflow。这就是前面说的“调度系统自己知道该跑哪天”的实现方式。但这里就引出了一个问题lastValue的值从哪里来如果你用的是lastmodified增量每次跑的last-value必须是从上次跑完得到的最新值。靠手动改job.properties不是自动化该干的事。实用做法是搞一个水位表专门存每张表的lastValue。在workflow里加一个前导的shell action先从水位表查出上次的值然后Sqoop任务用它做过滤跑完再更新水位表。这样整个链路才是真·自动化。示例的逻辑大致是# 在shell action里执行类似这样的逻辑 last_value$(hive -e SELECT last_value FROM sync_watermark WHERE table_nameorders) # 把last_value写到HDFS或者通过arg传给sqoop action sqoop import ... --last-value $last_value ... # 跑完后更新 hive -e UPDATE sync_watermark SET last_value... WHERE table_nameorders不过这种shell里调hive的做法在Oozie里嵌套比较麻烦我自己常用的是两个Hive action夹一个Sqoop action第一个Hive action读水位输出到一个HDFS临时文件第二个Sqoop action用--options-file读取这个临时文件里的last-value第三个Hive action再更新水位。稍微绕但稳。4.4 参数配置与资源控制job.properties里还需要注意YARN队列和资源设置。把调度任务固定到一个专门的队列比如dataSync里避免和白天实时任务抢资源。同时给每个map限制内存queueNameetl mapreduce.map.memory.mb2048 mapreduce.reduce.memory.mb2048 mapreduce.map.java.opts-Xmx1536m mapreduce.reduce.java.opts-Xmx1536m mapreduce.job.reduce.slowstart.completedmaps0.8这套配置意味着每个map最大吃2G内存JVM堆只分1.5G剩下的留给其他开销。如果你的集群节点内存不大比如单台32G那并发太高很可能直接让节点OOM——这个要在工具层面控制住不然半夜调度一跑第二天早上集群就被压垮了。还有一个实际经验Sqoop导入的map数量跟源库能承受的连接数直接相关。MySQL如果默认max_connections是150你一张表开20个map再加上其他业务在跑大概率会把你那台MySQL连接数打满。我一般做法是大表--m 8到--m 12小表--m 2就够宁可慢一点也别把源库压垮——源库一挂影响的不只是你的同步任务。5. 常见问题与排查技巧实录5.1 Sqoop连不上MySQL的排查标准路径这个真的是所有Sqoop新手都会遇到的第一个拦路虎。“sqoop连接不上mysql”相关的报错五花八门但梳理下来无非是下面几类排查顺序也基本固定先看驱动。报No suitable driver found for jdbc:mysql://...那十有八九是mysql-connector-java.jar没放对位置。Sqoop会扫描$SQOOP_HOME/lib目录下的所有jar需要把jar丢到这个目录后重启会话如果是oozie action还要确认workflow的lib里有。再看URL参数。报Communications link failure或者Connection refused这通常是网络不通。先telnet 数据库IP 3306测一下端口通不通通了再看是不是权限问题——用同一个账号在命令行里直接mysql连一下看能不能登上。如果是权限问题报错里会提到Access denied for user。这里有个细节MySQL的账号授权里host部分如果写的是localhost那你从别的机器连会直接被拒绝需要授权成%或者具体IP段。再看SSL和时区相关的报错。MySQL 8.0以上如果连串里没关SSL会报javax.net.ssl.SSLHandshakeException时区不对会报The server time zone value Öйú±ê׼ʱ¼ä is unrecognized——这种乱码错误信息就是时区问题的典型特征直接加上serverTimezoneAsia/Shanghai解决。有一个真实案例可以分享一下。有个朋友跑Sqoop报Connection refused但他本机用MySQL客户端连同一个库完全没问题。排查到一半发现他是在本机装了个Sqoop但Sqoop把作业提交到了YARN集群上而集群的DataNode没法访问数据库那台机器的IP——因为他只在本地hosts里写了映射。这个问题在不少公司都有解决方法是把源库的内网域名或者IP映射同步到所有集群节点上或者配置统一的内网DNS。5.2 Oozie作业运行失败后的日志追踪方法Oozie报错最头疼的地方在于它不是直接给你一个清晰的错误信息而是各种嵌套的日志让人一头雾水。但其实掌握了正确的日志查找顺序大多数问题三五分钟内就能定位。第一步先用oozie job -info查看作业状态oozie job -oozie http://oozie-server:11000/oozie -info 0000001-240615135040949-oozie-oe-B-1看结果里的Status字段。如果是RUNNING说明作业卡在某一步KILLED和FAILED都表示已经挂了。挂掉之后oozie job -log能拉出来这个作业的完整执行日志搜索关键词ERROR或者Exception基本上能看出是哪个环节、哪个action出的问题。第二步如果是Sqoop action挂了去YARN上看对应的MapReduce作业日志。yarn logs -applicationId application_xxx可以拿到Container级别的完整stderr和stdout。这里最可能出现的情况是Sqoop进程起来之后被YARN杀了说明分配的内存不够更加要关注mapreduce.map.memory.mb的配置。第三步如果Oozie本身报错去Oozie服务端日志目录翻oozie.log。在CDH里这个路径一般是/var/log/oozie/oozie.log通过grep job_id能找到更底层的异常堆栈。这里再分享一个超实用的技巧给Oozie作业配SLA。在coordinator.xml里加一段SLA定义就能在作业快超时的时候提前触发告警info sla:info sla:nominal-time${coordinator:nominalTime()}/sla:nominal-time sla:should-start5/sla:should-start sla:should-end360/sla:should-end sla:messageMySQL到Hive的同步作业超时/sla:message /sla:info /info这里should-start是期望作业从开始到启动的时间上限分钟should-end是期望总时长。Oozie会实时对比超出阈值就把告警推到配置的webhook或者邮件接口。这比干等着看告警群里的失败消息要优雅得多——尤其是半夜你可以在作业卡住不敢自己走完的时候第一时间就被paged到。5.3 增量数据正确性与数据质量校验自动化调度中最怕的不是任务失败而是任务“成功”了但数据是错的。增量场景里最典型的现象就是重跑和漏数据。我曾经历过一次比较惨痛的事故Sqoop的lastmodified模式下因为上游某条记录的更新时间早于last-value的边界被当成历史数据处理就丢在目标表里了。为了及时发现这种问题我在同步链路里加了一个校验步骤每次都让Sqoop统计本次导入的行数和源库按where条件查出来的行数比较。具体做法是在workflow里加一个shell action跑到最后执行一段比较逻辑# 源库行数 src_cnt$(mysql -h${jdbcHost} -u${jdbcUser} -p${jdbcPassword} -N -e SELECT COUNT(*) FROM orders WHERE update_time ${lastValue}) # 目标表分区行数 dst_cnt$(hive -e SELECT COUNT(*) FROM ods.orders WHERE dt${syncDate}) if [ $src_cnt -ne $dst_cnt ]; then echo count mismatch: source$src_cnt dest$dst_cnt # 发送告警不阻断下游任务但标记为需人工核查 exit 0 fi这里有个细节到底该用exit 1阻断后续任务还是exit 0只发告警取决于你的业务容忍度。如果是财务相关的表坚决阻断并重试如果是日志分析之类的可以只发告警不让链路卡死。我这边的经验是数据一致性没有灰色地带但影响范围可以分级——宁可让下游等一等也不要让脏数据进去。5.4 调度积压与集群负载的观察手段如果你跑了一段时间自动化调度早晚会遇到一个问题任务积压。Workflow明明是按天跑但某一天因为源库慢整个同步晚了3个小时下一天又因为上游数据量暴涨晚了5个小时叠加下去任务执行时刻越推越晚。这个时候光看单次任务成功还是失败已经不够了得看调度节奏本身的健康状况。我这边常用的几个观察点一是看Coordinator的每个nominal time理论执行时间和actual time实际执行时间之间的差距。Oozie的job info里会显示这两个时间如果差值越拉越大说明系统在持续积压。处理手段要么是扩容资源要么是拆任务——把大表拆成多个小任务并行或者把不重要的任务挪到非高峰时段跑。二是看调度间隔和作业时长的比例。如果某个作业平均要跑40分钟而调度频率是每小时一次长期看必定积压。这种场景要么降低频率要么优化Sqoop并发和Hive SQL没有别的捷径。三是看YARN集群的CPU和内存水位。在集群的ResourceManager页面上能看到每个队列的使用情况重点关注dataSync队列。如果队列里长期有任务排队等待说明资源已经到瓶颈了。这时候加map数没有意义——YARN不会因为你map多了就多给你资源只会让更多的任务排更久的队。我踩过这个坑以为是Sqoop并发不够狂调--m参数结果单个任务没变快反而让整个队列更堵。6. 从自动化到稳定的最后一公里6.1 告警和重试机制怎么设计才不算过度很多人在搭调度系统的时候会陷入一个误区把重试参数全拉满失败就一直重跑以为这就是稳定。实际上盲目重试往往会让问题变得更糟。第一条命令因为源库连接超时挂了重试5次后面4次大概率也是超时——因为这些重试都在同一时间发起踩的是同一个坑。我的经验是区分错误类型做决策连接类错误Connection refused、SocketTimeout、Communications link failure可以重试1到2次间隔10分钟。如果网络抖动是偶发这种重试能救回来但如果是源库彻底挂了超过2次重试就该果断告警别再往集群里堆垃圾作业了。数据类错误字段类型转换失败、主键冲突、NULL约束违规重试没有意义直接跳到告警分支发详细的错误日志给人看。资源类错误YARN队列满、磁盘空间不足不要立刻重试等到非高峰时段再拉起或者降低并发重新尝试。Oozie本身支持重试配置在action里加上retry-max2/retry-max retry-interval10/retry-interval这里的retry-interval单位是分钟。但我会额外留个心眼告警策略覆盖的workflow必须带error分支到告警action否则Oozie默认的失败处理就是静默退出你必须盯着oozie job的状态才能发现问题这等于没有自动化监控。6.2 备份、回滚和灰度上线的顶层设计自动化程度越高出问题时的影响范围就越大。所以我还要强调一个大家容易忽略的话题自动化调度系统本身也要有灾备措施。最基础的措施是定期备份Oozie的workflow目录和coordinator配置。我的做法是每周打包一次/user/bigdata/oozie/worksflows到单独的备份目录保留最近4周。出问题时快速回滚到上一版配置比现场改XML重新上传要快得多。再者上线新表或者修改同步逻辑不要直接改生产环境的coordinator先去测试环境跑通再拷贝到生产。这个说起来容易但我见过太多人直接在线上改配置然后马上-run结果半路挂了搞到整个队列的作业全部失败。稳妥做法是先在测试环境跑一遍确认无误后复制到生产目录然后重新提交coordinator任务。增量水位表也要有备份。真遇到极端情况——比如Hive里的水位表被误删你需要能重建水位。我一般会在MySQL里另存一张镜像表每周同步一次水位备份这样即使Hive这边乱了也能从MySQL恢复水位继续跑不至于整个增量链路瘫痪。6.3 从手工运维向平台化运营演进当你的同步任务从几张表变成几十张表、几百张表的时候靠一个个coordinator任务管理是不够的。我见过不少团队的管理方式就是拷贝几十份workflow目录每个表一个——然后维护成本直线上升改一个公共参数得几十个地方同步改漏掉任何一个就跑出坑。我自己的思路是把公共参数收敛到一份job.properties模板里每个表只维护差异化的部分表名、check-column、split-by、目标分区路径然后通过脚本批量生成workflow和coordinator配置。参数变更时只要改模板重新生成和提交即可。能用程序解决一遍的事儿别用手动做五十遍。更进一步可以把水位管理、任务监控、失败重跑这些逻辑都收口到一个管理平台上。比如开发一个小工具能批量查询所有coordinator的状态、批量重跑失败作业、批量更新表的水位。这一步做完才敢说自己把自动化数据迁移做成了体系。最后想说点实在的这套Sqoop加Oozie的调度系统我已经跑了快两年前后迭代过好几轮。要说最深的体会反而是最开始以为最难的“把数据导过去”其实是最简单的部分后面遇到的麻烦基本都在数据一致性、任务依赖、资源控制和告警恢复这些“边缘问题”上。如果你刚开始搭我建议第一步先别急着铺开几十个任务选一张核心业务表把链路完整跑通——从Sqoop导入、Oozie调度、到告警和重跑机制全部补齐确认的是这套方法论本身是可靠的再往里面加表。这个节奏看起来慢实际上反而是最快的路径因为每一类问题你都是在单张表上先踩完、先总结完再批量复制推广的。还有一个小经验Sqoop任务的连接串和Oozie的workflow XML都放在Git里管理。改配置留痕出问题一键回滚。这个习惯帮我躲过了好几次灾——别人在现场改得一团糟的时候我直接把上一版拉下来重新部署就完事了。自动化调度系统的价值不在于你写了几段好看的流程编排而在于半夜两点该跑的任务都在正常跑跑完了该给下游的东西一分不少真出了岔子告警能在第一时间找到人。这套东西做到了你就能睡个踏实觉了。
返回列表