Flink SQL与Table API在实时数据处理中的混合编程实践与调优
1. 项目概述当Flink SQL遇上复杂业务逻辑最近在做一个实时数据处理的POC项目核心需求是把来自Kafka的订单流和MySQL里的商品维表关联起来做实时统计和异常检测。团队里既有熟悉SQL的数据分析师也有擅长Java API的实时计算工程师。为了平衡开发效率和执行性能我们决定深度使用Flink的Table API SQL。这已经是我们在生产环境应用的第三个核心案例了前两个分别解决了CDC入湖和实时聚合的问题。这次面临的挑战更综合多流关联、自定义函数、状态调优以及端到端的一致性保证。通过这个案例我想分享的不只是几行SQL或API调用而是如何将Flink Table API SQL这套声明式的编程模型真正落地到解决复杂、多变的实时业务场景中并规避那些只有踩过坑才知道的陷阱。2. 整体架构与核心设计思路2.1 业务场景与需求拆解我们的业务场景是一个电商实时风控与大盘。数据源主要有两个一个是Kafka里面流淌着JSON格式的订单事件流包含order_id,user_id,product_id,amount,event_time等字段另一个是MySQL数据库存储着商品维表包含product_id,product_name,category,price等字段。需求可以拆解为三层实时关联与丰富需要将订单流实时与商品维表进行关联Lookup Join补全商品名称和类目信息形成一条完整的订单明细宽表。复杂指标计算基于宽表需要计算一些复杂指标例如同一用户10分钟内的下单频率防刷单、同一类目下的异常金额订单可能为标错价、热门商品排行榜等。这些计算涉及窗口、聚合和自定义的逻辑判断。多路输出处理结果需要写入多个目的地实时异常订单需要告警并写入Elasticsearch供风控系统查询统计结果如每分钟各类目GMV需要写入ClickHouse做实时大盘展示完整的明细数据还需要备份到HDFS或Iceberg表仓中供下游分析。2.2 技术选型为什么是Table API SQL面对这样的需求我们评估了三种实现方式纯DataStream API、纯Flink SQL、以及混合模式Table API SQL。最终选择了混合模式原因如下开发效率对于常规的流关联、过滤、分组聚合SQL的声明式写法极其高效。数据分析师可以直接参与SQL编写降低了沟通成本。例如一个简单的每分钟GMV统计一句SQL就能搞定远比用DataStream API手动开窗、聚合要快。灵活性当遇到SQL无法直接表达的复杂逻辑时比如需要访问状态进行自定义的模式检测我们可以无缝切换到Table API甚至最终转换成DataStream使用ProcessFunction。这种“可退可进”的能力是纯SQL或纯API难以比拟的。生态与维护Table API SQL在连接器Connector、格式Format、Catalog方面生态日益完善。通过SQL DDL来定义源表、维表和结果表结构清晰易于管理和版本化。后续新增一个输出目的地往往只需要加一段CREATE TABLE语句。性能优化Flink SQL引擎在底层做了大量优化如查询优化器、代码生成等。对于大多数标准操作其生成的执行计划往往比手写的通用DataStream作业更高效。我们只需要在少数热点处进行针对性调优。基于以上考虑我们的核心设计思路是用SQL处理80%的标准ETL和聚合逻辑用Table API和UDF处理15%的复杂业务逻辑最后用DataStream API处理剩下5%需要精细控制状态和时间的特殊场景。2.3 环境准备与依赖配置首先我们需要搭建一个可用的开发环境。这里以Flink 1.16版本和Java 11为例。Maven核心依赖properties flink.version1.16.0/flink.version /properties dependencies !-- Flink Table API SQL 基础依赖 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-table-api-java-bridge/artifactId version${flink.version}/version scopeprovided/scope /dependency !-- 计划器用于本地执行和优化 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-table-planner_2.12/artifactId version${flink.version}/version scopeprovided/scope /dependency !-- 如果需要与DataStream API互转需要此依赖 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java/artifactId version${flink.version}/version scopeprovided/scope /dependency !-- 连接器依赖 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-files/artifactId version${flink.version}/version /dependency !-- JSON格式解析 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-json/artifactId version${flink.version}/version /dependency !-- MySQL JDBC驱动 -- dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version8.0.33/version /dependency /dependencies注意在Flink 1.15及以上版本社区推荐在集群环境中使用flink-table-api-java和flink-table-runtime而将flink-table-planner标记为provided或排除以支持更灵活的编译。但在本地IDE开发调试时我们通常仍需要引入planner。生产打包时需要根据集群情况调整避免jar包冲突。3. 核心实现从表定义到复杂查询3.1 定义数据源与维表一切始于表的定义。我们使用StreamTableEnvironment的executeSql方法来执行DDL语句。1. 定义Kafka订单源表CREATE TABLE order_events ( order_id STRING, user_id BIGINT, product_id BIGINT, amount DECIMAL(10, 2), event_time TIMESTAMP(3), process_time AS PROCTIME(), -- 处理时间属性用于Lookup Join WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND -- 事件时间属性与水印 ) WITH ( connector kafka, topic order_topic, properties.bootstrap.servers kafka-broker:9092, properties.group.id flink-sql-demo, scan.startup.mode latest-offset, format json, json.fail-on-missing-field false, json.ignore-parse-errors true );这里有几个关键点PROCTIME()声明了一个处理时间属性这对于无界的维表JoinLookup Join是必须的因为Lookup Join目前只支持处理时间语义。WATERMARK定义了事件时间属性和水印生成策略。TIMESTAMP(3)表示毫秒精度的时间戳。水印延迟5秒用于容忍一定程度的数据乱序。这对于基于事件时间的窗口聚合至关重要。json.ignore-parse-errors设置为true是个好习惯可以防止因个别脏数据导致整个作业失败脏数据会被过滤掉并记录日志。2. 定义MySQL商品维表CREATE TABLE product_dim ( product_id BIGINT, product_name STRING, category STRING, price DECIMAL(10, 2), update_time TIMESTAMP, PRIMARY KEY (product_id) NOT ENFORCED -- 声明主键对JDBC连接器优化很重要 ) WITH ( connector jdbc, url jdbc:mysql://mysql-host:3306/demo_db, table-name product_info, username your_user, password your_password, lookup.cache.max-rows 10000, -- 缓存最大行数 lookup.cache.ttl 10min -- 缓存过期时间 );对于维表lookup.cache配置是性能关键。它将维表数据缓存在TaskManager内存中避免每条流数据都去查询数据库。max-rows和ttl需要根据维表数据量和更新频率权衡。我们的商品表大约1万条且变更不频繁所以设置10分钟缓存是合理的。如果维表更新频繁可以设置较小的TTL或使用lookup.async true进行异步查询但这会增加复杂度。3.2 实现流与维表的关联Lookup Join有了源表和维表关联就变得非常简单。我们使用LATERAL TABLE语法进行Lookup Join。-- 创建订单明细宽表视图 CREATE VIEW order_detail AS SELECT o.order_id, o.user_id, o.product_id, p.product_name, p.category, o.amount, p.price, o.amount / p.price as quantity, -- 计算数量 o.event_time FROM order_events AS o LEFT JOIN product_dim FOR SYSTEM_TIME AS OF o.process_time AS p ON o.product_id p.product_id;这段SQL看起来和批处理的Join没什么两样但它在流处理引擎下是持续运行的。FOR SYSTEM_TIME AS OF o.process_time指定了基于处理时间的时态表关联。左表order_events的每一行数据到达时都会根据当时的process_time去查询product_dim维表的最新快照进行关联。实操心得Lookup Join的维表必须声明主键PRIMARY KEY ... NOT ENFORCED否则Flink优化器可能无法生成最高效的执行计划。NOT ENFORCED表示Flink不会在运行时检查主键约束只是给优化器一个提示。3.3 复杂指标计算窗口聚合与自定义函数接下来我们基于order_detail视图计算复杂指标。1. 计算每分钟各商品类目的GMV事件时间窗口CREATE VIEW category_gmv_per_minute AS SELECT category, window_start, window_end, SUM(amount) as total_gmv, COUNT(DISTINCT order_id) as order_count FROM TABLE( TUMBLE(TABLE order_detail, DESCRIPTOR(event_time), INTERVAL 1 MINUTE) ) GROUP BY category, window_start, window_end;这里使用了TUMBLE窗口函数基于event_time字段划分1分钟的滚动窗口。window_start和window_end是窗口表值函数TVF产生的元数据字段非常方便。2. 检测用户高频下单滑动窗口与Having子句为了防止刷单我们需要检测10分钟内下单超过5次的用户。CREATE VIEW potential_risk_users AS SELECT user_id, window_start, window_end, COUNT(order_id) as order_cnt FROM TABLE( HOP(TABLE order_detail, DESCRIPTOR(event_time), INTERVAL 1 MINUTE, INTERVAL 10 MINUTE) ) GROUP BY user_id, window_start, window_end HAVING COUNT(order_id) 5;这里使用了HOP滑动窗口窗口大小10分钟滑动步长1分钟。这意味着每分钟都会输出过去10分钟内的统计结果能更及时地发现异常。3. 使用自定义函数UDF判断异常订单有些业务规则SQL写起来很别扭比如判断一个订单金额是否远高于该商品历史平均价格例如3倍标准差以外。这时就需要UDF。 首先用Java或Scala编写一个标量函数ScalarFunctionimport org.apache.flink.table.functions.ScalarFunction; import java.math.BigDecimal; public class IsAbnormalOrderFunc extends ScalarFunction { // 判断当前金额是否异常 public Boolean eval(BigDecimal currentAmount, BigDecimal avgPrice, BigDecimal stdDevPrice) { if (avgPrice null || stdDevPrice null) { return false; } BigDecimal threshold avgPrice.add(stdDevPrice.multiply(new BigDecimal(3))); return currentAmount.compareTo(threshold) 0; } }然后在TableEnvironment中注册这个函数tableEnv.createTemporarySystemFunction(IS_ABNORMAL, new IsAbnormalOrderFunc());最后在SQL中调用-- 假设我们已有商品的历史均价和标准差视图 product_stats CREATE VIEW abnormal_orders AS SELECT od.*, ps.avg_price, ps.stddev_price, IS_ABNORMAL(od.amount, ps.avg_price, ps.stddev_price) as is_abnormal FROM order_detail od JOIN product_stats ps ON od.product_id ps.product_id WHERE IS_ABNORMAL(od.amount, ps.avg_price, ps.stddev_price) TRUE;注意事项UDF虽然灵活但会阻碍Flink的部分优化如谓词下推且序列化/反序列化开销较大。应优先使用内置函数仅在必要时使用UDF并确保UDF函数是确定性的同一输入总是相同输出。3.4 多路输出Multiple Sinks计算结果需要写入多个目的地。我们为每个输出目标定义一个结果表然后使用INSERT INTO语句。1. 定义输出表-- 1. 异常订单写入Elasticsearch CREATE TABLE es_risk_orders ( order_id STRING, user_id BIGINT, amount DECIMAL(10,2), event_time TIMESTAMP(3), reason STRING, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector elasticsearch-7, hosts http://es-host:9200, index risk_orders, document-id.key-delimiter $ ); -- 2. 类目GMV写入ClickHouse CREATE TABLE ch_category_gmv ( category STRING, window_start TIMESTAMP(3), window_end TIMESTAMP(3), total_gmv DECIMAL(15,2), order_count BIGINT ) WITH ( connector jdbc, url jdbc:clickhouse://ch-host:8123/demo, table-name category_gmv, username default, password ); -- 3. 订单明细写入FileSystem (Parquet格式用于归档或Iceberg) CREATE TABLE fs_order_detail ( order_id STRING, user_id BIGINT, ... event_time TIMESTAMP(3) ) PARTITIONED BY (dt) WITH ( connector filesystem, path hdfs://namenode:9000/data/order_detail, format parquet, sink.partition-commit.delay 1 h, sink.partition-commit.policy.kind success-file );2. 执行插入理论上我们可以为每个INSERT INTO启动一个独立的作业。但更高效的方式是使用StatementSet将多个INSERT语句打包成一个作业执行共享源端计算减少资源消耗。TableEnvironment tableEnv ...; StatementSet stmtSet tableEnv.createStatementSet(); stmtSet.addInsertSql(INSERT INTO es_risk_orders SELECT order_id, user_id, amount, event_time, high_frequency as reason FROM potential_risk_users); stmtSet.addInsertSql(INSERT INTO ch_category_gmv SELECT * FROM category_gmv_per_minute); stmtSet.addInsertSql(INSERT INTO fs_order_detail SELECT *, DATE_FORMAT(event_time, ‘yyyy-MM-dd’) as dt FROM order_detail); // 执行所有插入 stmtSet.execute();这样一个Flink作业就会同时向三个目的地输出数据计算逻辑只在源头执行一次。4. 高级特性与性能调优实战4.1 状态管理与TTL配置流处理作业的核心是状态。上述作业中窗口聚合、Lookup Join缓存都会产生状态。如果状态无限增长最终会导致作业失败。我们必须为状态设置生存时间TTL。在Flink SQL中可以通过table.exec.state.ttl参数来设置空闲状态的保留时间。但更精细的控制需要在StreamTableEnvironment的配置中设置。StreamTableEnvironment tableEnv StreamTableEnvironment.create(env); Configuration configuration tableEnv.getConfig().getConfiguration(); // 设置空闲状态保留时间为10分钟 configuration.setString(table.exec.state.ttl, 10min);这个配置意味着如果一个Key如user_id在10分钟内没有接收到新的数据其对应的状态将被自动清理。这对于用户行为分析类任务非常有用可以自动清理不活跃用户的状态。踩坑记录TTL设置得太短可能导致晚到的数据无法与之前的状态正确关联例如一个用户9分钟没下单第10分钟下单之前的状态已被清理窗口计数会从0开始。设置得太长状态会膨胀。需要根据业务数据的实际时间分布来调整。对于事件时间窗口水印的推进也会触发过期状态的清理。4.2 并行度与资源优化并行度设置直接影响作业的吞吐量和延迟。我们的作业包含多个算子Source、维表Join、窗口聚合、Sink。Source并行度通常与Kafka Topic的分区数一致这是消费吞吐量的上限。维表Join并行度Lookup Join算子LookupJoinRunner是无状态的缓存是每个任务实例本地缓存的其并行度可以设置得较高以分散查询压力。但要注意每个并行子任务都会维护一个独立的缓存总缓存量 并行度 *lookup.cache.max-rows。窗口聚合并行度聚合算子WindowAggregate是有状态的其并行度决定了状态被分到多少个KeyGroup中。增加并行度可以分散热点Key的压力但也会增加状态快照Checkpoint的网络开销。一般建议设置为资源槽Slot数量的整数倍。Sink并行度需要根据下游系统的写入能力来定。例如写入ClickHouse如果ClickHouse表是分布式表可以设置较高的并行度如果是单点并行度设置过高可能导致写入冲突或连接数过多。可以在SQL中通过/* OPTIONS(‘parallelism’‘4’) */提示来为单个操作设置并行度但这属于实验性功能。更通用的方法是在提交作业时通过-p参数设置全局并行度或者在代码中env.setParallelism()。资源调优建议TaskManager堆内存状态大的作业需要更多堆内存。可以通过taskmanager.memory.process.size设置。托管内存Managed MemoryFlink用于排序、哈希表、缓存的状态后端如RocksDB会使用托管内存。对于状态很大的作业需要增加taskmanager.memory.managed.size或taskmanager.memory.managed.fraction。网络缓冲区高吞吐作业需要更多的网络缓冲区来避免背压。可以调整taskmanager.memory.network.min/max/fraction。4.3 利用CDC源实现维表实时更新我们之前用的JDBC维表缓存有TTL这意味着维表变更有延迟。如果商品价格变更需要实时反映在订单计算中比如计算实时的毛利率就需要使用CDCChange Data Capture源作为维表。Flink提供了flink-connector-mysql-cdc等工具。将商品维表定义改为CDC源CREATE TABLE product_dim_cdc ( product_id BIGINT, product_name STRING, category STRING, price DECIMAL(10, 2), update_time TIMESTAMP, PRIMARY KEY (product_id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname mysql-host, port 3306, username your_user, password your_password, database-name demo_db, table-name product_info, server-time-zone Asia/Shanghai );此时product_dim_cdc就是一个动态变化的维表。当MySQL中的商品信息发生INSERT、UPDATE、DELETE时这个变化会实时同步到Flink作业中。与之关联的order_detail视图也会实时产出更新后的结果。这实现了真正的实时数仓“拉链表”效果但代价是状态变得更大需要存储维表的所有版本且对源数据库的变更捕获有性能要求。5. 常见问题排查与调试技巧在实际开发和运维中我们遇到了不少问题这里总结几个典型的。5.1 “Cannot generate a valid execution plan for the given query”这是一个非常常见的错误通常意味着SQL语句在逻辑上或与配置结合时有问题。可能原因1Lookup Join的维表缺少主键声明。这是最容易被忽略的一点。务必在维表DDL中加上PRIMARY KEY (your_pk) NOT ENFORCED。可能原因2在批处理模式下使用了流表特有的函数比如PROCTIME()。检查执行环境是StreamTableEnvironment还是BatchTableEnvironment。可能原因3连接器或格式不支持所需的操作。例如某些版本的JDBC连接器作为维表时可能不支持带复杂条件的Join。查看对应连接器的官方文档。排查方法使用tableEnv.explainSql(sql)打印执行计划观察在哪个节点失败。简化SQL逐步添加JOIN、WHERE等子句定位问题点。5.2 数据延迟或背压Backpressure作业运行一段时间后发现输出延迟越来越高。可能原因1Sink写入慢。这是最常见的原因。检查目标系统如ClickHouse、ES的负载、写入速率是否达到瓶颈。可以观察Flink UI的Backpressure监控看压力是否集中在Sink算子。可能原因2维表查询慢。Lookup Join如果缓存未命中会同步查询外部数据库。如果数据库响应慢会成为瓶颈。可以查看数据库监控并考虑优化增大缓存、使用异步查询、对维表做预聚合或分片。可能原因3数据倾斜。某个商品类目或用户特别活跃导致聚合算子的某个子任务处理的数据量远大于其他子任务。在Flink UI的Metrics中查看每个算子的numRecordsInPerSecond是否均衡。解决方案对于Sink慢增加Sink并行度、批量写入、调整目标库的配置如ES的refresh_interval。对于维表查询慢优化数据库查询加索引、使用更快的缓存如Guava Cache升级为Caffeine、引入Redis作为二级缓存。对于数据倾斜在SQL层面尝试先对热点Key加随机后缀打散进行一轮聚合然后再去重汇总。或者使用/* SKEW(‘source’, ‘key’) */提示实验性功能。5.3 状态持续增长导致Checkpoint失败作业运行几天后Checkpoint开始超时或失败。可能原因1状态TTL未设置或设置过长。这是根本原因。为所有聚合和Join操作评估合理的数据有效期并设置table.exec.state.ttl。可能原因2RocksDB状态后端配置不当。如果使用RocksDB其LSM树结构在 compaction 时可能产生临时空间膨胀。确保taskmanager.memory.managed.size足够并考虑调整RocksDB的写缓冲区、压缩策略等高级参数。可能原因3有“僵尸”Key。某些Key可能不再有数据流入但其状态因为TTL未到期而一直存在。如果Key空间无限如用户ID这会导致状态无限增长。需要从业务上考虑是否能用有限的分组如按城市、按等级来替代无限分组。排查方法使用Flink的State Processor API或从Checkpoint元数据中分析状态大小分布。监控numRegisteredKeyedState和stateSize指标。5.4 SQL调试与日志排查调试Flink SQL作业不像调试Java代码那样直观。技巧1使用临时视图打印中间结果。将复杂的SQL拆解把中间步骤创建为TEMPORARY VIEW然后SELECT * FROM view_name并输出到Print或BlackHole连接器在本地IDE运行观察数据。技巧2开启Plan日志。在log4j.properties中设置log4j.logger.org.apache.flink.table.plannerDEBUG可以在日志中看到详细的逻辑计划和优化后的执行计划有助于理解Flink是如何翻译你的SQL的。技巧3利用Metric系统。Flink提供了丰富的Metrics如numRecordsIn、numRecordsOut、currentInputWatermark、stateSize等。在Flink UI或对接的监控系统如PrometheusGrafana中配置关键指标看板是线上运维的必备手段。特别要关注lastCheckpointDuration和lastCheckpointSize它们直接反映了状态健康度。6. 从Table API回退到DataStream的混合编程尽管Table API SQL强大但总有边界。当我们需要实现一个非常复杂的、基于状态的业务逻辑或者需要精确控制时间、处理迟到数据到侧输出流时回退到DataStream API是更佳选择。假设我们需要在风控中实现一个自定义的“序列模式检测”检测用户“浏览-加入购物车-下单”这个序列在1小时内是否异常频繁。用SQL实现非常困难但用DataStream的PatternProcessFunction则很合适。我们可以这样做将Table转换为DataStream// 假设 order_detail_with_action 表包含了用户行为类型view, cart, order Table orderDetailTable tableEnv.sqlQuery(SELECT user_id, action, event_time FROM order_detail_with_action); DataStreamRow orderDetailStream tableEnv.toDataStream(orderDetailTable); // 转换为POJO类型方便后续处理 DataStreamUserAction actionStream orderDetailStream .map(row - new UserAction(...)) .returns(UserAction.class);在DataStream上应用复杂模式检测PatternUserAction, ? riskyPattern Pattern.UserActionbegin(view) .where(new SimpleConditionUserAction() {...}) // 浏览事件 .next(cart).within(Time.hours(1)) .where(...) // 加购事件 .next(order).within(Time.hours(1)) .where(...); // 下单事件 PatternStreamUserAction patternStream CEP.pattern( actionStream.keyBy(UserAction::getUserId), riskyPattern ); DataStreamAlert alerts patternStream.process(new MyPatternProcessFunction());将检测结果再写回Table生态// 将DataStream转换回Table Table alertTable tableEnv.fromDataStream(alerts); // 注册为临时视图供后续SQL使用 tableEnv.createTemporaryView(risk_alerts, alertTable); // 可以继续用SQL将告警与其他数据关联或输出 tableEnv.executeSql(INSERT INTO es_alerts SELECT * FROM risk_alerts);这种混合模式充分发挥了两种API的优势SQL用于高效的数据准备和标准计算DataStream用于实现核心的、定制化的复杂业务逻辑。关键在于Table和DataStream之间的无缝转换这得益于Flink Table模块底层统一的RowData数据结构。7. 生产部署与监控告警开发调试完成后最终要部署到生产环境。我们使用Flink on Kubernetes的模式通过Flink Kubernetes Operator或Flink Session Cluster进行部署。作业提交与配置管理我们将所有的DDL和DML SQL语句维护在一个或多个.sql文件中。通过sql-client或编程方式读取并执行。# 使用sql-client提交 ./bin/sql-client.sh -f /path/to/job.sql # 编程方式 String sql Files.readString(Paths.get(“/path/to/job.sql”), StandardCharsets.UTF_8); tableEnv.executeSql(sql);将数据库连接信息、Kafka地址等敏感配置提取到配置文件中通过-D参数或环境变量注入避免硬编码。监控告警体系Flink作业本身通过REST API或与Prometheus集成监控numRestarts重启次数、lastCheckpointDuration上次Checkpoint耗时、numberOfFailedCheckpoints失败的Checkpoint数。任何非零的失败Checkpoint或频繁重启都需要立即告警。业务指标在作业中通过MetricGroup自定义业务指标如ordersProcessed、abnormalOrdersDetected。将这些指标暴露给Prometheus在Grafana中制作业务大盘。数据质量监控在关键的数据出口如写入ClickHouse的表设置数据量的波动监控同比、环比。如果某个时间点数据量骤降或为零可能意味着上游作业挂了或数据流异常。端到端延迟在源头Kafka消息中打入时间戳在Sink处比较当前时间与消息时间戳的差值作为端到端延迟指标进行监控。经过这个案例的实践我们团队对Flink Table API SQL的应用从“能用”到了“敢用”和“会用”的阶段。它确实极大地提升了开发实时数据处理的效率但同时也要求开发者对流处理的核心概念时间、状态、容错有深刻理解才能用好并调优。记住没有银弹在享受声明式编程便利的同时也要时刻关注作业的运行状态和资源消耗做好混合编程的准备以应对那些SQL无法覆盖的、千变万化的业务需求。