尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Flink sql应用范围

Flink sql应用范围 可以用几个实际场景来理解标准关系运算适合 Flink SQL复杂状态和流程控制适合 DataStream API。例子一SQL 可以完成的 CDC 同步需求把 MySQL 的订单表实时同步到 Iceberg。MySQL orders | v Flink CDC | v Iceberg orders如果只是字段映射、过滤和类型转换Flink SQL 就足够CREATETABLEmysql_orders(order_idBIGINT,user_idBIGINT,statusSTRING,amountDECIMAL(18,2),update_timeTIMESTAMP(3),PRIMARYKEY(order_id)NOTENFORCED)WITH(connectormysql-cdc,hostnamemysql-host,port3306,usernameflink,passwordpassword,database-nametrade_db,table-nameorders,scan.startup.modeinitial);CREATETABLEiceberg_orders(order_idBIGINT,user_idBIGINT,statusSTRING,amountDECIMAL(18,2),update_timeTIMESTAMP(3),PRIMARYKEY(order_id)NOTENFORCED)WITH(connectoriceberg,catalog-namehadoop_catalog,catalog-typehadoop,warehouses3://warehouse/);INSERTINTOiceberg_ordersSELECTorder_id,user_id,status,amount,update_timeFROMmysql_ordersWHEREamount0;这里的逻辑可以清楚地表达为输入表 - 过滤 - 字段选择 - 输出表因此使用 SQL 可读性更好开发和维护成本也较低。例子二SQL 可以完成窗口统计需求统计每个账户 5 分钟内的交易金额和交易笔数。SELECTaccount_id,window_start,window_end,COUNT(*)AStrade_count,SUM(amount)AStotal_amountFROMTABLE(TUMBLE(TABLEtrades,DESCRIPTOR(trade_time),INTERVAL5MINUTES))GROUPBYaccount_id,window_start,window_end;这种场景主要是按 Key 分组按事件时间开窗口做 Count、Sum 等聚合。这是 Flink SQL 的典型适用场景。例子三简单去重适合 SQL需求同一订单只保留最新的一条记录。SELECTorder_id,user_id,status,update_timeFROM(SELECT*,ROW_NUMBER()OVER(PARTITIONBYorder_idORDERBYupdate_timeDESC)ASrnFROMorder_events)WHERErn1;如果“最新”的定义就是按照一个可靠的更新时间或版本号排序SQL 可以直接完成。但如果版本判断需要结合多个来源的日志位点、事务 ID 和特殊业务规则就可能需要 DataStream API。例子四复杂状态机不适合只用 SQL需求监控账户是否出现以下交易模式5 分钟内发生 3 次失败交易 - 随后 1 分钟内发生一笔大额交易 - 触发风险告警这个逻辑涉及保存每个账户的失败交易列表失败记录超过 5 分钟后清理失败次数达到 3 次后进入中间状态注册 1 分钟定时器收到大额交易时触发告警超时后清理状态并恢复初始状态。使用KeyedProcessFunction会更直观publicclassRiskProcessFunctionextendsKeyedProcessFunctionLong,Trade,RiskAlert{privateValueStateIntegerfailedCount;privateValueStateLongexpireTimer;OverridepublicvoidprocessElement(Tradetrade,Contextctx,CollectorRiskAlertout)throwsException{IntegercountfailedCount.value();if(countnull){count0;}if(FAILED.equals(trade.getStatus())){count;failedCount.update(count);if(count3){longtimerctx.timestamp()60_000L;ctx.timerService().registerEventTimeTimer(timer);expireTimer.update(timer);}}elseif(count3trade.getAmount().compareTo(newBigDecimal(100000))0){out.collect(newRiskAlert(trade.getAccountId(),trade.getTradeId()));}}OverridepublicvoidonTimer(longtimestamp,OnTimerContextctx,CollectorRiskAlertout)throwsException{failedCount.clear();expireTimer.clear();}}Flink SQL 可以表达其中一部分过滤、窗口和聚合但要把这种多阶段状态转换完整地写成 SQL通常会变得复杂且难以维护。例子五需要动态规则时使用 Broadcast State需求风控规则不是写死在程序里的而是从 Kafka 动态下发规则 1账户 5 分钟失败次数 3 规则 2单笔金额 100000 规则 3高风险地区禁止交易交易流需要和规则流关联交易流 ------------------ -- 动态风控处理 规则流 - Broadcast State-规则更新后所有并行实例都要及时获得新规则。这个过程通常使用BroadcastStreamBroadcastProcessFunction广播状态自定义规则匹配逻辑。SQL 可以做普通维表 Join但对于动态广播规则、复杂规则版本切换和自定义匹配过程DataStream API 更合适。例子六外部异步调用适合 DataStream API需求每笔交易都需要查询外部风控服务交易事件 - 调用风控 HTTP 服务 - 根据返回结果决定放行或拦截如果直接在 SQL UDF 中同步调用 HTTP 服务可能导致算子线程阻塞外部服务变慢时产生反压难以控制并发数超时和重试逻辑不清晰Checkpoint 延迟增加。更适合使用AsyncFunctionpublicclassRiskAsyncFunctionextendsRichAsyncFunctionTrade,EnrichedTrade{OverridepublicvoidasyncInvoke(Tradetrade,ResultFutureEnrichedTraderesultFuture){riskClient.checkAsync(trade.getAccountId()).thenAccept(result-{resultFuture.complete(Collections.singletonList(newEnrichedTrade(trade,result)));}).exceptionally(error-{resultFuture.completeExceptionally(error);returnnull;});}}这样可以显式控制最大并发请求数超时时间重试次数失败降级异步结果和原事件的关联。例子七自定义 Sink 不适合只用 SQL需求将结果写入一个内部交易系统该系统要求每批最多 500 条 请求失败重试 3 次 使用 trade_id version 幂等 连续失败触发熔断 失败消息写入补偿队列Flink SQL 只能表达“把结果写入某个 Sink”但这些工程行为通常需要自定义 Sink 或 DataStream 算子实现结果流 - 批量聚合 - 异步发送 - 重试 - 熔断 - 失败补偿 - 幂等提交例子八自定义协议数据源不适合 SQL需求公司内部设备通过 TCP 长连接发送二进制报文报文头 | 设备 ID | 时间戳 | 类型 | Payload | 校验码需要自己完成TCP 连接管理半包和粘包处理二进制解码校验码验证断线重连位点保存错误报文旁路输出。这种场景要实现自定义 Source 或 DataStream 算子Flink SQL 本身不能直接完成协议接入。例子九两者混合使用很多生产任务不会只选一种 API而是让 SQL 和 DataStream API 各自负责擅长的部分Flink CDC Source | v DataStream API | 复杂去重、状态机、异步风控 v 转换为 Table | v Flink SQL | Join、窗口聚合、字段加工 v Iceberg / Doris例如DataStreamTradetrades...;DataStreamRiskResultresultstrades.keyBy(Trade::getAccountId).process(newRiskProcessFunction());tableEnv.createTemporaryView(risk_results,results);tableEnv.executeSql( INSERT INTO risk_summary SELECT account_id, COUNT(*), SUM(amount) FROM risk_results GROUP BY account_id );一个简单判断标准能否清楚地描述为输入表 - 过滤/Join/聚合 - 输出表 能优先 Flink SQL 不能需要自定义状态、定时器、外部调用或协议处理DataStream API 两者都有前半段或后半段使用 SQL中间复杂逻辑使用 DataStream API可以概括为例如 Kafka 到 Iceberg 的字段转换、窗口聚合和普通 CDC 同步可以用 Flink SQL但账户风控状态机、动态规则广播、异步调用外部服务、自定义 TCP Source、带重试和补偿的 Sink则需要 DataStream API 或自定义 Connector。实际项目中通常混合使用两者。
返回列表