流计算能力功能说明

更新时间:
复制 MD 格式

云消息队列 Kafka 版内置流计算能力,兼容标准流式 SQL 语法,可直接使用 SQL 表达实时处理逻辑,无需编写代码,即可在 Kafka 内完成流式数据的读取、计算与写出。

标准流式 SQL

流计算能力提供完整的流式 SQL 支持,开发者可用熟悉的 SQL 语法定义实时处理逻辑。核心能力如下。

基本 SQL 语法

能力

说明

数据定义(DDL)

创建、修改、删除表、视图、函数等。

数据查询(DQL)

支持 SELECTFROMWHEREGROUP BYHAVINGORDER BYLIMITOFFSET;支持子查询(标量/行/表子查询)、JOIN(INNER/LEFT/RIGHT/FULL OUTER/CROSS/LATERAL)、窗口函数、表达式计算、类型转换、CASE WHEN 等。

数据操作(DML)

INSERT INTO 将查询结果写入目标表(支持追加或 UPSERT 模式);流模式下不支持 UPDATE/DELETE 语句。

事务性 DML

依赖 Sink 实现幂等或事务性写入(如 Kafka 事务)。

流处理能力

能力

说明

时间属性

支持事件时间(Event Time)与处理时间(Processing Time)。

水印(Watermark)

处理乱序事件,控制延迟容忍度。

窗口(Windowing)

支持滚动窗口(Tumbling)、滑动窗口(Hopping)、会话窗口(Session)、累积窗口(Cumulative)、基于元素数量的窗口(Count-based)。

窗口 TVF

支持 TUMBLEHOPSESSIONCUMULATE

内置函数

类别

示例

标量函数

UPPER()SUBSTRING()ROUND()COALESCE()

聚合函数

COUNTSUMAVGMAXMINARRAY_AGGJSON_OBJECTAGG

窗口函数

ROW_NUMBER()RANK()LEAD()LAG()FIRST_VALUE()

时间函数

CURRENT_TIMESTAMPDATE_FORMATTIMESTAMPDIFFLOCALTIMESTAMP

JSON 函数

JSON_VALUEJSON_QUERYJSON_EXISTSJSON_OBJECTJSON_ARRAY

条件函数

CASE WHENIFNULLIFCOALESCE

类型转换

CAST(value AS TYPE)

数据格式

支持 JSONCSVAvro(支持 Schema Registry)、ParquetORCDebezium-JSON(CDC 变更日志)、Canal-JSONRaw(原始字节)等多种格式的解析与写入。

连接器生态

流计算能力内置连接器,覆盖上下游主流系统的读写与同步,可在一个流任务内完成"读取—处理—写出"的完整链路。

类别

支持系统

消息队列

Kafka、MQTT、RocketMQ

说明

连接器会逐渐开放,如有紧急需求,可通过工单提出。

云原生弹性与免运维

  • 自动伸缩:系统根据流任务的实时负载自动扩展和缩减计算资源,兼顾性能与成本。

  • 免集群运维:无需自建和维护流处理集群,无需关注与 Kafka 的版本兼容性。