如果有人让你在一套数仓里同时扛住实时写入、离线ETL、即席查询和报表输出,你会怎么设计?这就是我最近在做的Flink与Greenplum集成项目:用Flink承担实时计算和增量数据管道,用Greenplum接住大规模并行分析,两者互相配合应对典型的混合负载。这篇文章是把整个项目里踩过的坑、选型的权衡、落地的步骤完整梳理一遍,给正在做实时数仓、HTAP、或者被“又要实时又要跑大查询”折磨的人做个参考。
1. 为什么要把Flink和Greenplum放一起:混合负载的场景拆解
1.1 混合负载到底指什么,真正难在哪
混合负载这个词看起来抽象,落到实际业务里非常具体:白天业务高峰,实时订单数据每分钟到达,需要立刻清洗、聚合、写入;同一时间,运营同学在跑日报查询,数据分析师在刷大屏看趋势,甚至还有几个定时ETL任务在重算历史指标。这些任务共享同一套基础设施,看着像“读写并发”,实际比读写并发要复杂得多。
我习惯打个比方:这就像一个餐厅厨房,既要同时处理外卖快餐的短平快订单,又要准备晚宴的几十道菜。快餐讲究快,晚宴讲究稳和出品一致。如果所有订单都挤在同一口锅里,谁都快不了。混合负载最麻烦的地方就是资源无法简单静态分配,因为流任务的写入是持续的、不打招呼的,而分析查询往往是突发性的、吃大资源的,两者一旦碰上,经常是查询慢、写入也堆积。
另外一个难点是一致性。流计算天然是“边到边”的处理模型,数据可能重复、可能乱序,而分析系统希望看到的是干净、稳定的结果。所以混合负载不是一个组件能搞定的,需要一个流处理引擎负责动态数据的加工,一个MPP分析引擎负责沉底的数据存储和复杂查询。这也是Flink和Greenplum组合的价值所在。
1.2 Flink和Greenplum各自的位置
Flink在实时计算里属于“全能型选手”。它能做无界流处理,也能做有界批量处理;支持事件时间、窗口计算、状态管理、Checkpoint和端到端的Exactly Once语义。我们通常把它放在数据管道的中游,负责从Kafka、CDC或各种业务库把数据捞出来,做实时清洗、关联、聚合,然后发给下游存储。
Greenplum则完全另一种脾气。它是基于PostgreSQL的MPP架构数据库,底层是shared-nothing分布式的Segment节点,擅长把一张大表拆到多个Segment上并行扫描。配合appendonly列存表,跑几亿行的聚合统计、复杂Join,比传统单机数据库快得多。它的定位就是“沉底分析”:把已经加工好的明细或汇总数据放进来,对外提供BI报表、多维分析和数据挖掘。
这两个东西不是替代关系,而是上下游关系。Flink负责“算”,Greenplum负责“存和查”。把Flink算完的结果丢到Greenplum,再由Greenplum接住高并发查询,这是混合负载里最典型的分工。
1.3 集成后的目标形态和能解决的业务问题
我这次做的项目,最核心的场景有三个:
第一个是实时用户行为分析。埋点日志从Kafka进来,Flink按用户维度做分钟级聚合,算出PV、UV、转化率,写入Greenplum的汇总表。业务方打开大屏,看到的是延迟不到一分钟的实时指标。
第二个是实时风控。交易事件进入Flink后,需要在毫秒级关联账户维度和历史交易特征,这个维表就挂在Greenplum上。Flink实时拉取GP里的最新维度信息做Join,判断这笔交易是否异常。
第三个是T+0报表。以前跑一份全量报表要等到凌晨批量算,现在Flink持续把当天增量数据写进GP,分析人员随时可以查当日累计数据,晚上再用批任务处理历史归档。
这三个场景覆盖了“实时写入、在线分析、维表关联”三类负载,放在一起才叫真正的混合负载。如果只做一个维度,比如只是Flink到GP的写入,那集成价值会小很多。
2. 集成设计绕不开的四个关键决策
2.1 通道选择:直连、缓冲还是批量文件
Flink和Greenplum之间怎么传数据,决定了整个链路的吞吐、延迟和运维复杂度。我实际对比过三种主流方式,也分别试过坑。
第一种是Flink JDBC连接器直接写Greenplum。配置最简单,写SQL就行,适合每秒几千条的小流量场景。缺点是Greenplum的写入路径偏OLAP,高并发的小事务写入容易碰到锁冲突,而且JDBC分批提交的语义并不完美,数据量一大性能就会出现波动。这个方案的延迟是秒级,但吞吐天花板很低。
第二种是中间加一层Kafka。Flink把结果写到Kafka,再从一个独立任务消费Kafka批量写入GP。这样做的好处是削峰填谷,Flink不用关心Greenplum是不是慢,由Kafka缓冲住瞬时流量。缺点是链路变长,多了一套组件,延迟也从秒级变成了十几秒到分钟级。但换来的是稳定性和可扩展性。
第三种是文件落地加Greenplum外部表加载。Flink算完以后写Parquet到HDFS或对象存储,Greenplum用gpfdist或PXF外部表去读。这种方式吞吐极高,适合日级或小时级的大批量数据同步,但实时性最差,不适合增量报表。
我最终的取舍是:实时增量走Kafka到Flink,再由Flink批量写GP;离线补数和历史初始化走外部表加载。直连JDBC只用在测试环境或极低流量场景。如果你在项目里拿不定主意,先想想你的QPS、可接受的延迟、运维人力这三件事,就能选出自己那条路。
2.2 表模型与分布键设计
Greenplum不是单机数据库,表数据是按分布键散到各个Segment的。分布键选得不好,就会发生数据倾斜:某个Segment塞满了数据,其他Segment空闲,查询和写入都受累。很多团队把Flink的数据直接塞进GP,却没设计表分布,结果慢得怀疑人生。
和Flink集成时,表模型设计的核心逻辑是:GP表的分布键必须和Flink写入的数据特征对齐。比如Flink按user_id做分组聚合,那么GP表也应该用user_id做分布键。这样不管是Flink写入还是后续GP查询按用户维度聚合,数据都能在本地完成大部分计算,而不是跨Segment广播。
除此之外,一定要考虑分区表。Greenplum的分区和分布是两个独立概念,分布决定数据落在哪个Segment,分区决定数据在Segment里怎么切片。我建议按时间做Range分区,比如每天一个分区,这样Flink写入当天分区,历史分区可以转换为只读,避免和增量写入抢资源。如果业务要保留90天数据,还可以直接drop掉90天前的分区,比delete快得多。
还有一点容易被忽略:如果GP表是列存表(AO_COLUMN),非常适合分析查询,但高频的小批量写入性能不如行存表。实时明细表建议用行存或堆表,指标汇总表用列存。要根据访问模式分开设计。
2.3 一致性:从at-least-once到最终幂等
Flink的Checkpoint机制可以保证作业重启后不会丢数据,但默认的JDBC Sink不支持真正的Exactly Once。任务重启后,可能有一部分数据已经写进Greenplum,Flink又会重放一次,导致重复。如果你对数据准确性要求高,必须在下游设计幂等写入。
最简单有效的办法是给目标表设置业务主键,然后利用数据库的upsert能力。需要说明的是,Greenplum的版本差异很大,7.x基于PostgreSQL 12,原生支持ON CONFLICT;6.x还是老版本,不支持这个语法。如果碰到6.x,常规做法是写一个自定义JDBC Sink,先UPDATE再INSERT,或者用一个临时表承接增量,再通过GP的merge语句合并。
我在项目里更推荐一种稳妥做法:Flink写入GP时,把数据打上批次ID和时间戳,目标表保留一个批次字段。万一发生重复,可以按批次字段快速定位并清理。这套“标记+清理”的思路比单纯依赖数据库事务要可靠得多,尤其当Greenplum参与混合负载时,不可能为了一个流任务开长事务。
还有一点必须提醒:Flink的Checkpoint间隔不要太短,否则整个链路频繁对齐状态,反而影响吞吐。一般业务场景5到10分钟一个Checkpoint,配合GP侧幂等清理,最终一致性完全够用。
2.4 资源隔离:写和查不能互相拖死
如果Flink直接往Greenplum灌数据,同时BI查询也在跑,会出现一个典型现象:大查询占用大量Segment内存和CPU,导致Flink的写入事务迟迟无法提交,然后JDBC连接堆积,写入延迟升高。反过来,如果实时写入一直占着资源,分析查询又被拖慢。这是混合负载设计里最容易被低估的一环。
Greenplum支持资源队列和资源组两种资源管理方式。资源组(Resource Group)更灵活,可以限制CPU使用率、内存占用和并发数。我给Flink写入任务单独建了一个资源组,限制并发数为5到10,CPU上限控制在20%左右,避免它和BI查询抢资源。
Flink侧的隔离也要做。如果你的Flink跑在YARN上,就给实时任务单独划分一个队列,配置独立的内存和CPU;如果是K8s,就通过namespace或者ResourceQuota隔离。不要把所有Flink任务混在一个默认队列里,特别是那些每小时跑一次的批任务,很容易干扰实时任务。
另外要控制Greenplum查询侧的并发,尤其是大查询。很多BI工具会同时发出大量并发SQL,即使单条SQL不慢,并发一高也能把GP拖垮。我在前面加了一层SQL限流,限制大查询并发最多3个。这部分效果,比调任何参数都明显。
3. 核心实操:从零搭一条Flink+Greenplum实时分析链路
3.1 版本与环境准备
这套东西版本搭配很关键,我最初用的是Flink 1.14 + Greenplum 6,后来升到Flink 1.18 + Greenplum 7。升级最大的好处是Greenplum 7支持了ON CONFLICT和更好的事务处理,JDBC写入的幂等方案简单了很多。
环境清单如下:
- Flink 1.18,运行在YARN集群上,TaskManager内存按作业配置。
- Greenplum 7,Master节点和8个Segment节点。
- Kafka 3.x,作为流数据缓冲层。
- PostgreSQL JDBC驱动42.x,Flink官方Jdbc Connector依赖自带驱动或手动上传。
连接Greenplum的JDBC URL和普通PostgreSQL很接近:
jdbc:postgresql://gp-master:5432/analytics注意Greenplum默认端口是5432,如果Master有多个网卡,要确认Flink集群能连通Master的对外地址。连接用户名建议用专用于同步的账号,不要直接用超级管理员,这样后续在GP资源组里隔离更方便。
3.2 Flink SQL实现流式聚合写Greenplum
我这次用的是Flink SQL,开发速度比DataStream API快很多,SQL提交就能跑。核心链路是Kafka源表、聚合查询、GP结果表三张表。
先建Kafka源表:
CREATE TABLE kafka_events ( user_id INT, event_type STRING, amount DECIMAL(10,2), event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND ) WITH ( 'connector' = 'kafka', 'topic' = 'user_events', 'properties.bootstrap.servers' = 'kafka-1:9092,kafka-2:9092', 'properties.group.id' = 'flink-gp-pipeline', 'format' = 'json', 'scan.startup.mode' = 'latest-offset' );再建GP结果表:
CREATE TABLE gp_user_agg ( user_id INT, event_count BIGINT, total_amount DECIMAL(12,2), window_start TIMESTAMP(3), PRIMARY KEY (user_id, window_start) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:postgresql://gp-master:5432/analytics', 'table-name' = 'user_agg', 'username' = 'etl_user', 'password' = '********', 'sink.buffer-flush.max-rows' = '1000', 'sink.buffer-flush.interval' = '5s' );最后执行INSERT:
INSERT INTO gp_user_agg SELECT user_id, COUNT(*) AS event_count, SUM(amount) AS total_amount, TUMBLE_START(event_time, INTERVAL '1' MINUTE) AS window_start FROM kafka_events GROUP BY TUMBLE(event_time, INTERVAL '1' MINUTE), user_id;提交任务后,Flink会每一分钟输出一个窗口结果,攒满1000条或者5秒周期触发写入。实际测试下来,这个配置能稳定支撑每秒几千条事件流。
有一个坑必须提醒:Flink官方JDBC Sink默认是append-only,即使你在WITH里声明了主键,它也不会帮你生成upsert语句。上面的SQL如果业务发生重复窗口重算,会插入重复数据。所以更稳妥的方式是使用DataStream API配合自定义SQL,写成INSERT INTO ... ON CONFLICT DO UPDATE。下面是一个简化示例:
JdbcExecutionOptions execOptions = JdbcExecutionOptions.builder() .withBatchSize(1000) .withBatchInterval(Duration.ofSeconds(5)) .build(); JdbcConnectionOptions connOptions = JdbcConnectionOptions.builder() .withUrl("jdbc:postgresql://gp-master:5432/analytics") .withDriverName("org.postgresql.Driver") .withUsername("etl_user") .withPassword("********") .build(); sink = JdbcSink.sink( "INSERT INTO user_agg(user_id, event_count, total_amount, window_start) " + "VALUES (?, ?, ?, ?) " + "ON CONFLICT (user_id, window_start) DO UPDATE SET " + "event_count = EXCLUDED.event_count, total_amount = EXCLUDED.total_amount", (ps, record) -> { ps.setInt(1, record.userId); ps.setLong(2, record.eventCount); ps.setBigDecimal(3, record.totalAmount); ps.setTimestamp(4, record.windowStart); }, execOptions, connOptions );Greenplum 7能直接跑这段SQL,6.x需要改写为UPDATE加INSERT的方式。自定义Sink写起来不复杂,但解决重复数据问题一劳永逸。
3.3 用Greenplum维表做实时关联
除了把Flink结果写进GP,很多时候也要反过来从GP读数据。最常见的场景是维表关联:实时流里的user_id只有ID,需要关联用户姓名、城市、等级。如果把维表放在GP里,Flink又需要实时获取,就可以用Flink SQL的维表Join。
先定义GP维表:
CREATE TABLE gp_dim_user ( user_id INT PRIMARY KEY NOT ENFORCED, user_name STRING, city STRING, level STRING ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:postgresql://gp-master:5432/analytics', 'table-name' = 'dim_user', 'username' = 'etl_user', 'password' = '********', 'lookup.cache.max-rows' = '5000', 'lookup.cache.ttl' = '10min' );主查询把实时流和维表关联:
SELECT e.user_id, d.user_name, d.city, e.amount FROM kafka_events AS e LEFT JOIN gp_dim_user FOR SYSTEM_TIME AS OF e.event_time AS d ON e.user_id = d.user_id;这个玩法让Greenplum从单纯的“分析仓库”变成了“可查询的数据服务”。注意维表Join是同步查询,每来一条数据都要访问GP,因此必须开缓存。我把缓存开到5000行、TTL 10分钟,既能拿到相对新的维度信息,又不会把GP Master压垮。如果维度更新特别频繁,可以把lookup.cache.ttl调低,但性能会下降,需要平衡。
3.4 资源配比和调优参数
这套链路搭建完,性能能不能起来,很大程度取决于参数怎么设。我总结出几个关键公式和原则,都是实测下来最有用的。
Greenplum写入并行度不要超过Segment数的两倍。比如8个Segment的集群,Flink的写入并行度建议8到16。超过这个值,Segment端的锁竞争和上下文切换会拖慢写入,看似并行高,性能反倒下降。
JDBC Sink的批次大小和间隔要配合。我的经验是batch-size >= 500、interval >= 3s。如果批次太小,提交事务太频繁,Greenplum的Master节点会成为瓶颈;如果批次太大,单条数据的端到端延迟会变高。1000条/5秒是一个通用的起点,再按流量微调。
Checkpoint间隔和状态大小也要考虑。短间隔(1分钟以内)会让实时任务频繁做快照,影响吞吐;长间隔(10分钟以上)又会让重启恢复变慢。我一般配置为3到5分钟,加上增量Checkpoint。
Greenplum侧最重要的调优是资源组。我给Flink写入账号设置了一个独立的资源组:
CREATE RESOURCE GROUP rg_flink WITH ( CONCURRENCY = 10, CPU_RATE_LIMIT = 20, MEMORY_LIMIT_PERCENT = 20 ); ALTER ROLE etl_user RESOURCE GROUP rg_flink;给BI查询账号设置另一个资源组:
CREATE RESOURCE GROUP rg_bi WITH ( CONCURRENCY = 5, CPU_RATE_LIMIT = 50, MEMORY_LIMIT_PERCENT = 50 ); ALTER ROLE bi_user RESOURCE GROUP rg_bi;这样Flink的写入最多占用20%的CPU,BI查询最多占用50%,两边都有明确上限。即便某一侧突然打满,也不会拖垮对方。再配合Hints或外部SQL限流,混合负载基本能稳定运行。
4. 排障实录:那些年踩过的坑
4.1 Flink JDBC连接器异常定位
我在项目里遇到最多的就是Flink写GP报连接异常,典型表现是作业刚启动就报Cannot connect to PostgreSQL server或者Connection is not available, request timed out。
先检查网络和驱动。Greenplum虽然兼容PG协议,但Master和各Segment之间还有内部通信,如果Flink作业跑在独立机房,到GP Master的网络延迟高,连接容易被Master端关闭。这时需要检查pg_hba.conf是否允许Flink节点的IP访问,以及密码认证方式。
第二个常见问题是线程池连接耗尽。Flink Sink默认使用HikariCP连接池,每个写入子任务都会建连接。如果并行度是16,连接池大小又不够,就会出现请求超时。解决方案是把jdbc.connection.max-retry-timeout调大,或者在JDBC URL里配合连接池配置增加上限。
第三个坑是驱动版本不匹配。Flink官方Jdbc Connector默认带驱动,但不同小版本支持的PostgreSQL协议有差异。Greenplum 6需要老一点的驱动,Greenplum 7用42.x就行。如果驱动版本太新或太旧,会报Protocol violation或者Unsupported authentication method。解决方法是把正确的驱动jar放到Flink的lib目录,并清掉容器里老版本的驱动。
4.2 写入性能上不去
Flink往Greenplum写数据,性能上不去有非常typical的几个原因。第一个就是目标表没有分布键或者分布键分布严重不均。比如表用distributed randomly,几个Segment数据量差异不明显,但写入时Master会随机分配,仍然容易产生网络开销。更好的做法是选高基数、业务查询Join的字段做分布键。
第二个原因是表上有太多索引。Greenplum的索引主要用于点查,分析型负载一般不建太多索引。索引多了,写入时每个Segment都要维护索引,性能下降很厉害。我建议流式写入的表只保留主键索引,或者干脆不用索引,通过分区裁剪来加速查询。
第三个原因是目标表的行存表使用了let's VACUUM不够。GP的堆表更新和删除会留下死元组,长期不清理会拖慢扫描。对高频写入的表,尽量用appendonly或定期执行VACUUM。
第四个原因特别容易被忽略:Flink写入批次设置过小。有人图延迟低,把buffer-flush.max-rows设为1,结果每个事务只插一条数据,Greenplum被频繁提交打爆。这种场景Low延迟反而是幻觉,因为GP的Master节点处理事务开销远高于单行插入的收益。
4.3 数据重复和丢数
混合负载链路上数据丢了或重复了,排查起来最费劲。丢数先看Flink的Checkpoint是否正常完成。如果Checkpoint一直失败,Kafka消费位点不会提交,重启后数据会重放,表现出来就是Greenplum里出现重复数据。
针对重复,我之前讲过用ON CONFLICT做幂等写入。这里再补充一种适用于Greenplum 6的兜底方案:在目标表上创建一个merge任务,每天凌晨把当天的临时表数据合并到主表。Flink只负责往临时表里写,合并任务负责去重和覆盖。虽然多了一步,但对老版本GP非常稳。
还有一种“看起来丢数”的情况:Flink的窗口聚合结果没有输出,是因为水位线没推进。Kafka源表的WATERMARK设置和事件时间必须匹配。如果业务时间比处理时间滞后很多,窗口会一直不触发,结果迟迟不写入GP。这时要检查Kafka消息的时间字段是不是标准格式,以及是否存在无界乱序。最简单暴力的方法是把watermark改为processing time,但这样就不是真正的流式语义了。
4.4 混合负载下查询变慢
在当前混合负载场景里,GP查询变慢往往不是因为SQL写得差,而是被实时写入资源抢占。我见过最多的情况是:Flink作业并行度拉满,持续写数据,BI客户端同时在跑大聚合查询,结果整个集群的Segment CPU全部打满,单个查询的返回时间从2秒变成40秒。
解决思路有三个层次。第一层,给Flink和BI划分不同的资源组,像前面那样限制CPU和并发。很多团队不做这步,本质上是让两个流量在同一个池子里抢资源,必然互相伤害。
第二层,在GP侧开启resource group的内存限制,配合statement_mem给查询合理分配内存。大查询如果申请不到足够内存,就会进入排队而不是直接卡死,反而保护了整体稳定性。
第三层,从时间维度错峰。如果业务不要求24小时实时分析,可以把Flink批量任务的写入时间放到凌晨或午饭低峰期。白天只跑轻量级增量写入,把重量级历史分析放到晚上。这不是技术妥协,而是混合负载架构设计里很实用的运营手段。
还有就是BI侧做查询治理。有些报表工具会发出没有谓词的SELECT *或者全表聚合,这种语句在GP里也能把Segment拖垮。我加了一层拦截规则,凡是扫描行数预估超过一定阈值的SQL,强制走队列,让它们排队执行而不是并发挤爆。
写在最后
这套Flink与Greenplum集成方案,我从选型、搭建到排障完整跑了大半年,最大的体会有三个:第一,混合负载不是一个纯技术问题,必须从业务流量、资源隔离、数据一致性三个维度同时设计;第二,Greenplum的并行能力很猛,但一定要把分布键、分区、资源组这些基本功打扎实,否则再强的MPP也白搭;第三,Flink到GP不管用什么通道,都要提前想好幂等策略,不然运维会让你天天处理重复数据。
最后分享一个我在实际项目里用出来的小技巧:给所有Greenplum表都保留一个etl_insert_time字段,Flink写入时带上系统时间。这样不但能排查数据延迟,还能在需要修复数据时,一条SQL按照时间范围精准删除或重算,不用对着整张表发愁。混合负载的链路越复杂,这种“留一手”的设计越值钱。