
前阵子帮一家银行做日终批量链路优化核心任务是把核心系统每天导出的TXT余额文件经过清洗和汇总之后落到PostgreSQL。这个流程拆开看并不复杂——无非是读文件、做聚合、写库但真正从零搭起来从解析规则到幂等控制再到Spark参数调优每一步都有不少细节。这篇文章就把整个“Spark TXT PostgreSQL”的余额处理流程完整复盘一遍包括环境搭建、文件解析、聚合逻辑、入库方式以及我在生产环境里踩过的坑。这篇适合谁看手里正拿着银行/财务系统TXT文件不知道怎么高效处理的工程师打算用Spark替换传统脚本做批量任务的团队以及任何想把Spark计算结果稳定写入PostgreSQL的同学。1. 场景理解与整体设计1.1 银行系统为什么还在用TXT文件交换数据银行业务系统对接尤其核心系统和外围系统之间最“古老”也最常见的交换格式就是TXT。你可能会问为什么不用接口、不用JSON现实是很多核心系统是十几年前的老架构对外提供文件接口最稳妥一根专线把文本文件丢过来上下游都不需要改造。而且日终批量本身就适合文件方式一次导出一整天的流水Spark凌晨定时扫描目录处理完入库链路简单可靠。这批TXT文件的格式也很有年代感竖线分隔或定长字段、GBK编码、偶尔带BOM头首行可能还有标题行。余额文件的典型字段包括账号、户名、币种、期初余额、借方发生额、贷方发生额、期末余额、更新时间。单文件可能几十万到上百万行全行账户多的时候日终批量要处理的数据量在千万级流水以上。这种量级用单机脚本跑确实能完成但性能和容错都不够稳。1.2 为什么选择Spark而不是Python脚本或ETL工具我第一次接到这个需求时第一反应也是“直接用Python读文件、处理、写库不就行了吗”但冷静看完数据量和业务要求之后就放弃了。原因有三千万级流水在单机上用Pandas处理内存压力很大而且出错了要整批重跑平衡文件可能有多个要做跨文件的关联和聚合脚本写起来容易失控银行场景要求可监控、可重试、可扩展以后数据量翻倍不能再推倒重来。Spark的优势在于把“分而治之”做到了框架层面。一个大文件默认会被切成多个分区每个分区由不同的Executor并行处理某个节点挂了会自动重试整体任务跑完还能看到每个Stage的耗时和Shuffle数据量。这也是Spark在银行ETL里被广泛使用的原因——处理能力强调度和监控也成熟。那目标库为什么选PostgreSQL一方面银行内部对Oracle/MSSQL的license费用敏感PostgreSQL开源且够用另一方面Spark自带PostgreSQL JDBC支持操作方便。PostgreSQL在数据一致性、复杂SQL、索引能力等方面也都够成熟。整个ETL的分层我按常见的“ODS-DWD-DWS”来设计ODS层保留TXT原始数据原样落地DWD层做清洗和规范化DWS层做账户余额汇总。很多团队对“etl的ods层”的概念有困惑这里明确一下ODS层就相当于数据仓库的“临时储物间”接进来什么就存什么先不纠结业务口径。算一下当时的整体数据量级单个账户汇总文件约80万行日交易流水文件约500万行两台Spark节点的集群8 core 32G每台跑完整条链路大概20分钟。这个表现在听起来不高但在生产环境里稳定性比绝对速度更重要。2. 环境准备与数据摸底2.1 Spark集群搭建与资源参数注意点我这边用的是Spark on YARN的部署方式。搭建环节不多说直接给一套在8核32G单机上验证过的核心参数模板。注意这里的参数不是越大越好而是要结合YARN队列的实际资源来设参数名推荐值说明spark.executor.instances4Executor数量spark.executor.cores2每个Executor占用CPU核数spark.executor.memory8gExecutor堆内内存spark.executor.memoryOverhead2g堆外内存用于JVM开销、Netty等spark.dynamicAllocation.enabledfalse先关掉动态分配便于稳定排查spark.sql.shuffle.partitions20Shuffle后默认分区数约cores*2~3很多人在Spark on YARN上都会遇到“CPU只能用1个”的问题我当时也被坑过。现象是YARN上每个Container只分配了1个vcore整个作业跑得极慢。原因通常是Spark默认的spark.executor.cores就是1如果你用spark-shell或默认提交方式启动且没有在提交脚本里显式指定核数YARN分配到的就是1个vcore。另一个隐藏坑是动态分配spark.dynamicAllocation.enabledtrue时如果下边界设置不当Executor可能始终维持最小个数。后面第6章会详细展开排查步骤。PostgreSQL这边安装时也容易翻车。常见报错是“无法创建锁文件 /var/run/postgresql/.s.pgsql.5432.lock: 权限不够”这多半是postgres用户对/var/run/postgresql目录没有写权限或者压根没初始化数据目录。解决方式也很简单重新mkdir并chown给postgres用户然后通过pg_ctl或pg_ctlcluster启动别直接用root硬启。目标表的设计我放第5章先看文件。2.2 看懂手上的TXT文件样例与摸底拿到文件后我习惯先做三件事查编码、看行格式、统计脏数据量。查编码可以用file命令file balance_20250601.txt # 输出balance_20250601.txt: ISO-8859 text, with very long lines这种ISO-8859通常是GBK/GB18030编码Spark读取时需要用encodinggbk或gb18030否则中文会变乱码。行格式则用head -3看一眼字段分隔符和字段顺序确认是否有表头行。我当时手上的TXT样例大概是这样的竖线分隔最后一行是汇总行ACCOUNT_NO|ACCOUNT_NAME|CURRENCY|BEGIN_BALANCE|DEBIT_AMOUNT|CREDIT_AMOUNT|END_BALANCE 6222001234567890123|张三|CNY|100.50|200.00|50.00|150.50 6222001234567890456|李四|CNY|200.00|0.00|100.00|300.00 ... TOTAL|1000000行|CNY|10000000.00|3000000.00|2000000.00|11000000.00摸底阶段还要确认文件末尾有没有汇总行、有没有空行、有没有非法金额。建议先跑一个简单的Spark count和格式校验把问题行数统计出来再决定解析方案。这个环节看似琐碎但能省掉后面不少麻烦——我见过有人因为TXT文件里有个肉眼不可见的BOM头导致第一列字段解析错位查了一下午才定位到。3. TXT文件解析与数据清洗3.1 Spark读取TXT文件的方式选择Spark读取纯文本文件有三类常用方式spark.read.text(path)每行作为一列value的DataFrame适合统一解析rdd.textFile(path)返回RDD[String]灵活但更底层spark.sparkContext.wholeTextFiles(path)每个文件作为一个record适合小文件但大文件会OOM不推荐。在ETL场景里我推荐用spark.read.text因为它可以直接利用DataFrame API做后续操作解析逻辑写在SQL或者withColumn里都方便。读取时注意设置编码和分区df_raw spark.read \ .option(encoding, gbk) \ .option(wholetext, false) \ .text(/data/etl/balance/20250601/balance_20250601.txt)这个DataFrame默认分区数取决于HDFS块大小和文件大小。如果只有一两个大文件分区数可能很少后续算子并行度不够。可以读取后显式repartition也可以设置spark.sql.files.maxPartitionBytes默认128MB来控制分区大小。不过更稳妥的做法是读进来之后按逻辑分区处理这个后面讲。3.2 从原始行到结构化DataFrame拿到value列之后按分隔符解析。竖线分隔的文件可以用split然后用when或cast做类型转换。我的解析逻辑大致是这样的from pyspark.sql.functions import split, trim, col, when, regexp_replace df_parsed df_raw.select( split(col(value), \\|).alias(arr) ).select( col(arr).getItem(0).cast(string).alias(account_no), col(arr).getItem(1).cast(string).alias(account_name), col(arr).getItem(2).cast(string).alias(currency), col(arr).getItem(3).cast(decimal(18,2)).alias(begin_balance), col(arr).getItem(4).cast(decimal(18,2)).alias(debit_amount), col(arr).getItem(5).cast(decimal(18,2)).alias(credit_amount), col(arr).getItem(6).cast(decimal(18,2)).alias(end_balance) )这里有个很关键的坑银行TXT里的金额字段可能是“100.50”这种正常格式但也可能是“1,000.50”带千分位甚至前面带正负号。我在解析前会统一做一次regexp_replace(col(value), ,, )再转decimal避免隐式类型转换失败导致整列变成null。另外如果源文件里有全角数字或中文括号也要先规范化否则cast后全是null。字段解析完后不要急着写库先做一轮基础校验账号长度是否为15~32位、金额是否大于等于0、日期字段是否合法。这些规则用Spark的filter和when都很容易实现。3.3 脏数据质检与异常隔离处理脏数据的关键原则是不要一有坏行就让整个任务挂掉也不要悄悄吞掉坏行。我的做法是双路输出正常数据进入后续加工异常数据写入独立的bad table/bad目录同时记录错误原因和原始行内容。Spark里实现“双路输出”可以用filter配合when标记错误类型。例如df_checked df_parsed.withColumn( is_valid, when(col(account_no).rlike(^[0-9]{15,32}$), 1).otherwise(0) ).withColumn( err_msg, when(col(is_valid) 0, account_no_invalid) .when(col(begin_balance).isNull(), begin_balance_null) .otherwise(ok) ) df_success df_checked.filter(col(is_valid) 1) df_failed df_checked.filter(col(is_valid) 0)异常数据单独落到PostgreSQL的etl_bad_record表字段包括source_file、line_no、raw_line、err_msg、etl_date。这样第二天早上看报表谁对不上能一眼定位。另外提一下编码。我遇到过上游系统导出的TXT是GB18030但文件头部却写着UTF-8结果Spark按UTF-8读出来一堆乱码。排查这类问题的手段是先把文件下载到本地用iconv或Python的chardet判断真实编码确认后再写死读取编码不要依赖文件名或环境变量。4. 余额处理核心逻辑4.1 业务口径期初、借贷发生、期末余额处理的核心口径简单说是期末余额 期初余额 贷方发生额 - 借方发生额但在实际银行项目里要注意科目的借贷方向。存款类科目对客户来说是贷方增加而资产类科目是借方增加。具体业务规则需要跟清算/核心系统对口径。我当时处理的是账户余额约定debit_amount表示借方发生额credit_amount表示贷方发生额对存款账户而言期末余额期初余额贷方发生额借方发生额。还有一点容易忽略同一个账号可能在TXT文件里出现多行比如一个多币种账户按币种拆行了所以不能直接按行累加要按账号币种维度做groupBy。这一步也正是Spark批处理的强项。4.2 按账号聚合reduceByKey 还是 groupByKey做余额汇总时我们需要把同一个账号的多行记录合并成一条汇总记录。Spakr RDD API里groupByKey把同一key的所有value拉到一个Iterator里reduceByKey则先在每个分区内做局部聚合再跨分区合并。显然reduceByKey的Shuffle量小得多大文件场景性能差距明显。所以优先用reduceByKey。用DataFrame API其实就是groupBy().agg()底层优化器会自动选择更高效的执行方案。核心代码from pyspark.sql import functions as F df_summary df_success.groupBy( col(account_no), col(currency) ).agg( F.first(account_name).alias(account_name), F.sum(begin_balance).alias(begin_balance), F.sum(debit_amount).alias(debit_amount), F.sum(credit_amount).alias(credit_amount), F.sum(end_balance).alias(end_balance) )然后按口径字段校验begin_balance credit_amount - debit_amount是否等于end_balance有差额的账号要单独标记出来。4.3 应对热门账号的数据倾斜余额处理里最典型的数据倾斜是极少数账号比如财政代发户、清算账户一天交易量巨大单key数据占了整个RDD的30%以上导致处理时其它Executor闲得慌这几个Executor忙到OOM。处理思路一般是加盐salt打散热点Key两阶段聚合第一阶段给聚合key拼接一个随机后缀比如0~9先做局部聚合第二阶段去掉后缀再对半聚合结果做最终聚合。这个方案能明显缓解倾斜但代价是需要保证同一个账号的所有记录不能因为加盐而跨批次丢失所以盐的范围和分配要均匀。更简单的替代方案是热点账号单独走一个高并行度子任务冷数据走常规链路。我当时用加盐方案效果还不错压测时整体耗时从23分钟降到16分钟。5. 写入PostgreSQL的工程化实现5.1 目标表设计与批次管理PostgreSQL目标表的设计重点不是字段多花哨而是如何支撑“每天重跑、幂等不重不漏”。我建的表结构类似这样CREATE TABLE account_balance_agg ( etl_date date NOT NULL, batch_id varchar(32) NOT NULL, account_no varchar(32) NOT NULL, currency varchar(8) NOT NULL, account_name varchar(128), begin_balance numeric(18,2), debit_amount numeric(18,2), credit_amount numeric(18,2), end_balance numeric(18,2), update_time timestamp DEFAULT now(), PRIMARY KEY (etl_date, batch_id, account_no, currency) );etl_date和batch_id是幂等控制的关键。每天跑批任务时先删除当日batch_id对应的数据再插入新数据。这样即使任务跑了一半失败重跑时不会产生重复记录业务方查数也永远看到的是“当天最新批次的完整结果”。5.2 Spark JDBC写入模式与参数调优Spark DataFrame写JDBC最常用的调用是df_summary.write \ .mode(append) \ .option(batchsize, 5000) \ .option(truncate, true) \ .option(numPartitions, 8) \ .jdbc(urljdbc:postgresql://10.10.10.10:5432/bank_etl, tableaccount_balance_agg, properties{user:etl_user, password:****, driver:org.postgresql.Driver})write的mode有四种append、overwrite、errorifexists、ignore。注意overwrite并不是真正的“覆盖”它会先drop表再创建如果你只想清空当日分区而保留表结构千万别用overwrite。正确做法是先执行DELETE FROM account_balance_agg WHERE etl_date ... AND batch_id ...再append写入。写入性能的关键参数batchsize默认1000每次批量插入的行数建议5000~10000太大容易撑爆PG端内存numPartitions控制写入并发度JDBC写入的并发不是越多越好建议2~4倍目标库CPU核数socketTimeout、connectTimeout网络抖动时避免连接卡死。还有很重要的一点在JDBC URL上拼接参数rewriteBatchedInsertstrue这个参数能让PostgreSQL把多条insert改写为批量插入实测写入速度能提升好几倍。如果不加这个参数Spark的batchsize设置对PG来说几乎无效每行还是一条单独insert慢到怀疑人生。5.3 幂等性与失败重跑写过批处理的人都知道最怕的不是任务失败而是任务失败后留下一堆“半成品”数据。Spark写JDBC天然不具备跨批次事务性所以要用业务手段保证幂等任务启动时先记录批次号处理完数据后在同一个批次事务里先删旧数据、再插新数据插入完成后更新批次状态表标记该日期该批次成功下次重跑时看到成功状态可以直接跳过或者覆盖重跑。删除旧数据时注意我一般不用TRUNCATE因为TRUNCATE会锁全表影响其它正在查询的业务。用DELETE WHERE etl_date ?配合索引能减少锁粒度。如果表很大可以考虑按天做分区表PostgreSQL原生支持range分区这样“删除当日”直接变成DROP PARTITION速度极快对在线查询影响也小。5.4 性能表现参考拿我当时的环境举例80万行余额汇总Spark集群2个节点、8核32G写入PostgreSQL4核16GnumPartitions8、batchsize5000、开启rewriteBatchedInsertstrue时写入耗时大约1~2分钟。如果不开启rewriteBatchedInserts同样的数据量跑完要20多分钟差距非常大。这个数据不一定适用于所有人的环境但可以作为线下压测的参考基线。6. 常见问题与排查实录6.1 Spark on YARN 只分配 1 个 CPU 核心这个问题的完整排查思路在YARN ResourceManager页面看Container的vcore数确认是不是全部为1检查Spark配置spark.executor.cores、spark.task.cpus默认1、spark.dynamicAllocation.enabled;检查YARN调度器的yarn.scheduler.maximum-allocation-vcores如果队列上限本身就小Executor申请再多核也拿不到检查提交命令是否用了--num-executors和--executor-cores如果你只设置了--num-executors没设cores默认cores就是1。解决模板提交时显式指定spark-submit \ --master yarn \ --executor-cores 2 \ --num-executors 4 \ --executor-memory 8g \ --conf spark.dynamicAllocation.enabledfalse \ ...我把这个配置套用到生产后作业的CPU利用率立刻上来了整体时间也缩短了一半以上。6.2 Executor OOM 问题余额处理里OOM多见于汇总/Join阶段。典型的异常是java.lang.OutOfMemoryError: Java heap space或ExecutorLostFailure。排查步骤查看Spark UI里哪个Stage的Shuffle Read/Write量最大检查单个分区数据量是否过大热点账号导致查看spark.sql.autoBroadcastJoinThreshold是否因为小表广播过大而撑爆Executor。调优方向增大spark.executor.memory和memoryOverhead把spark.sql.shuffle.partitions调大但要防止小文件过多对热点Key做加盐处理对大表Join时关闭自动广播或者手动指定broadcast join的小表。6.3 写入PostgreSQL很慢、锁等待写入慢通常是下面几个原因没有开启rewriteBatchedInserts每条insert都走单条SQLbatchsize设得太小网络往返次数太多目标表索引太多插入过程中索引维护开销巨大唯一索引冲突导致大量等待锁。我遇到过一次比较典型的因为表上有唯一约束(etl_date, batch_id, account_no, currency)但源数据里同一个key出现了重复行导致ON CONFLICT DO NOTHING没有生效大量写入卡在锁上。后来在写入前先用dropDuplicates去重问题才解决。去重代码很简单df_final df_summary.dropDuplicates([etl_date, batch_id, account_no, currency])6.4 中文乱码与字段错位这类问题属于“现象明显、根因隐蔽”的典型。乱码多半是编码识别错误字段错位多半是分隔符不一致或引号转义。建议在解析逻辑里加一层字段数量检查size(arr)不等于预期字段数时整行归入bad记录而不是强行取值。这样即使上游某天突然多了一个字段也不会悄无声息地把金额塞到账号列。6.5 快速排查速查表现象可能原因解决方案YARN每个Container只有1个vcoreexecutor.cores未设置或动态分配限制提交时显式指定--executor-cores关闭动态分配Executor OOMshuffle分区数少、热点数据、内存参数不足调大memory/memoryOverhead提高分区数加盐写入PG极慢未开启rewriteBatchedInsertsJDBC URL加rewriteBatchedInsertstrue写入锁等待主键冲突、索引多、mode误用先dropDuplicates设置合理batchsize用append配合DELETE中文乱码编码识别错误file命令chardet判断真实编码明确指定字段错位分隔符不一致、BOM头增加字段数校验清洗时去掉BOM整套流程跑下来我最想提醒的是开发阶段用一个小抽样文件反复验证解析规则生产阶段再全量跑。TXT解析、余额聚合、PostgreSQL写入这些环节单独看都不算难难的是把它们串成一条能每天稳定自动运行的链路。先想清楚文件格式、业务口径和重跑策略再动手写Spark代码比急着调参数值钱得多。如果你正好也在做类似的ETL任务希望这篇复盘能帮你少走几步弯路。