apache streampark 是一个专注于流处理应用程序开发和管理的开源平台,旨在简化 apache flink 和 apache spark 等流处理引擎的使用复杂度。它提供了一站式的应用开发、部署、监控和运维能力。
以下是关于 apache streampark 的功能介绍、使用场景及详细使用步骤示例。
一、 核心功能介绍
- 应用开发与管理
- 代码托管: 支持 git/svn 拉取代码,也支持直接在平台上编写 sql/scala/java 代码。
- 多引擎支持: 原生支持 apache flink 和 apache spark,同时兼容 datastream api、table api 和 sql。
- 版本管理: 内置应用版本控制,支持一键回滚。
- 自动化构建与部署
- 自动构建: 集成 maven/gradle/sbt,自动完成项目编译打包。
- 多种部署模式: 支持 remote、yarn-application、yarn-session、kubernetes session/application/native 等多种部署模式。
- ci/cd 集成: 提供 rest api,易于与 jenkins/gitlab ci 集成。
- 智能运维与监控
- 状态感知: 实时追踪作业状态(created, running, failed, finished 等)。
- 自动重启: 配置失败重试策略,异常时自动恢复。
- checkpoint/savepoint 管理: 可视化触发、查看和从指定 savepoint 恢复。
- 告警通知: 支持邮件、钉钉、企业微信、webhook 等多种告警渠道。
- sql 在线开发
- 提供 web ide 进行 flink sql 编写。
- 支持语法高亮、智能提示、sql 校验和血缘分析。
- 资源与配置管理
- 统一管理 flink home、maven 仓库、系统参数。
- 支持动态参数注入,区分开发/测试/生产环境配置。
二、 典型使用场景
| 场景 | 描述 | streampark 价值 |
|---|---|---|
| 实时数仓 etl | kafka → flink → doris/starrocks/clickhouse | 简化 sql 作业管理,统一调度数百个 etl 任务 |
| 实时监控大屏 | 业务指标实时聚合计算 | 快速部署、秒级状态反馈、故障自动恢复 |
| cdc 数据同步 | mysql/pg binlog → 下游存储 | 管理复杂的 cdc connector 配置和断点续传 |
| 算法模型推理 | 实时特征工程 + 模型预测 | 支持自定义 jar 包部署,管理模型文件版本 |
| 多租户平台 | 企业内部大数据平台 | 权限隔离、资源配额、统一入口降低使用门槛 |
三、 详细使用步骤示例(以 flink sql 作业为例)
前置准备
# 1. 克隆并启动 streampark (docker 方式最简) git clone https://github.com/apache/incubator-streampark.git cd incubator-streampark/docker docker-compose up -d # 2. 访问 web ui # 默认地址: http://localhost:10000 # 默认账号: admin / streampark
step 1: 配置 flink 环境
进入 setting → flink home,添加你的 flink 安装路径:
flink name: flink-1.17 flink home: /opt/flink-1.17.1 # 容器内或主机路径
step 2: 创建 flink sql 应用
进入 streampark → application → add new
| 配置项 | 示例值 | 说明 |
|---|---|---|
| development mode | custom code / flink sql | 选择 flink sql |
| execution mode | remote / yarn-application | 根据集群选择 |
| flink version | flink-1.17 | 关联已配置的 flink |
| application name | kafka-to-doris-demo | 应用名称 |
step 3: 编写 sql 代码
在 sql 编辑器中输入:
-- 源表定义
create table source_kafka (
user_id bigint,
item_id bigint,
behavior string,
ts timestamp(3),
watermark for ts as ts - interval '5' second
) with (
'connector' = 'kafka',
'topic' = 'user_behavior',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json',
'scan.startup.mode' = 'latest-offset'
);
-- 目标表定义
create table sink_doris (
window_start timestamp(3),
window_end timestamp(3),
behavior string,
cnt bigint
) with (
'connector' = 'doris',
'fenodes' = 'doris-fe:8030',
'table.identifier' = 'db.behavior_summary',
'username' = 'root',
'password' = ''
);
-- 核心逻辑
insert into sink_doris
select
tumble_start(ts, interval '1' minute) as window_start,
tumble_end(ts, interval '1' minute) as window_end,
behavior,
count(*) as cnt
from source_kafka
group by tumble(ts, interval '1' minute), behavior;step 4: 配置运行参数
# parallelism & checkpoint 配置
parallelism.default: 2
execution.checkpointing.interval: 30s
execution.checkpointing.mode: exactly_once
state.backend: hashmap
state.checkpoints.dir: hdfs:///streampark/checkpoints
state.savepoints.dir: hdfs:///streampark/savepoints
# 自定义参数(可在sql中通过 ${param} 引用)
kafka.brokers: kafka:9092step 5: 构建与启动
- 点击 build → 平台自动校验 sql 语法并解析依赖
- 点击 start → 选择是否从 savepoint 恢复
- 观察 flame graph / timeline 确认作业进入 running 状态
step 6: 运维操作
- 触发 savepoint: application 列表 → 操作栏 → savepoint → trigger
- 查看日志: 点击应用名 → log 标签页 → 实时查看 tm/jm 日志
- 修改 sql: stop (with savepoint) → edit sql → build → start (from savepoint)
- 告警配置: setting → alert → 添加钉钉/邮件组 → 绑定到应用
四、 最佳实践建议
- 生产环境务必开启 checkpoint,并配置
state.checkpoints.dir到 hdfs/s3 - stop 时选择 “with savepoint”,避免数据丢失
- sql 作业优先使用 catalog 管理表元数据,避免重复 ddl
- k8s 部署推荐使用 native 模式,获得更好的资源弹性
- 定期清理过期 checkpoint/savepoint,防止存储膨胀
- 利用 pipeline 功能 将多个关联 sql 合并为一个作业,减少资源开销
到此这篇关于apache streampark 功能和使用场景介绍、使用步骤详细示例的文章就介绍到这了,更多相关apache streampark 功能介绍内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
发表评论