加载内容...
署名-非商业性使用-禁止演绎
LOADING...
正在加载页面资源
LOADING...
正在初始化系统
1from pyflink.datastream import StreamExecutionEnvironment
2from pyflink.table import StreamTableEnvironment
3
4env = StreamExecutionEnvironment.get_execution_environment()
5t_env = StreamTableEnvironment.create(env)
6
7# 定义Kafka源表
8t_env.execute_sql("""
9 CREATE TABLE user_events (
10 user_id STRING,
11 event_type STRING,
12 event_time TIMESTAMP(3),
13 properties STRING
14 ) WITH (
15 'connector' = 'kafka',
16 'topic' = 'user-events',
17 'properties.bootstrap.servers' = 'kafka:9092',
18 'format' = 'json'
19 )
20""")
21
22# 实时聚合统计
23t_env.execute_sql("""
24 CREATE TABLE event_stats (
25 event_type STRING,
26 event_count BIGINT,
27 window_start TIMESTAMP(3)
28 ) WITH (
29 'connector' = 'clickhouse',
30 'url' = 'clickhouse://localhost:8123',
31 'table-name' = 'event_stats'
32 )
33""")
34
35t_env.execute_sql("""
36 INSERT INTO event_stats
37 SELECT
38 event_type,
39 COUNT(*) as event_count,
40 TUMBLE_START(event_time, INTERVAL '1' MINUTE) as window_start
41 FROM user_events
42 GROUP BY
43 event_type,
44 TUMBLE(event_time, INTERVAL '1' MINUTE)
45""")1-- ClickHouse 优化查询示例
2SELECT
3 toStartOfHour(event_time) as hour,
4 event_type,
5 count() as cnt,
6 uniqExact(user_id) as uv
7FROM user_events
8WHERE event_date >= today() - 7
9GROUP BY hour, event_type
10ORDER BY hour DESC
11LIMIT 1000| 场景 | 性能 | 说明 |
|---|---|---|
| 数据摄入 | 100,000 条/秒 | Kafka + Flink |
| 聚合查询 | <1 秒 | 亿级数据 |
| 明细查询 | <0.5 秒 | 百万级数据 |
| 存储压缩 | >5:1 | 列式存储 |