79人参与 • 2026-08-18 • ar
apache streampark 是一个专注于流处理应用程序开发和管理的开源平台,旨在简化 apache flink 和 apache spark 等流处理引擎的使用复杂度。它提供了一站式的应用开发、部署、监控和运维能力。
以下是关于 apache streampark 的功能介绍、使用场景及详细使用步骤示例。
| 场景 | 描述 | streampark 价值 |
|---|---|---|
| 实时数仓 etl | kafka → flink → doris/starrocks/clickhouse | 简化 sql 作业管理,统一调度数百个 etl 任务 |
| 实时监控大屏 | 业务指标实时聚合计算 | 快速部署、秒级状态反馈、故障自动恢复 |
| cdc 数据同步 | mysql/pg binlog → 下游存储 | 管理复杂的 cdc connector 配置和断点续传 |
| 算法模型推理 | 实时特征工程 + 模型预测 | 支持自定义 jar 包部署,管理模型文件版本 |
| 多租户平台 | 企业内部大数据平台 | 权限隔离、资源配额、统一入口降低使用门槛 |
# 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
进入 setting → flink home,添加你的 flink 安装路径:
flink name: flink-1.17 flink home: /opt/flink-1.17.1 # 容器内或主机路径
进入 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 | 应用名称 |
在 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;# 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:9092state.checkpoints.dir 到 hdfs/s3到此这篇关于apache streampark 功能和使用场景介绍、使用步骤详细示例的文章就介绍到这了,更多相关apache streampark 功能介绍内容请搜索代码网以前的文章或继续浏览下面的相关文章希望大家以后多多支持代码网!
您想发表意见!!点此发布评论
版权声明:本文内容由互联网用户贡献,该文观点仅代表作者本人。本站仅提供信息存储服务,不拥有所有权,不承担相关法律责任。 如发现本站有涉嫌抄袭侵权/违法违规的内容, 请发送邮件至 2386932994@qq.com 举报,一经查实将立刻删除。
发表评论