
整理了一份Grafana项目实战案例教程:实时交易热力图。 为有相关需求的技术人员提供一个思路和一个技术指导。希望能给大家带来帮助。
为金融、电商等高频交易场景提供实时可视化解决方案
维度 | 指标示例 | 可视化形式 |
|---|---|---|
时间 | 每秒交易量 | 时间轴热力图 |
地理 | 城市/经纬度交易密度 | 地理热力图 |
业务维度 | 商品类目/价格区间分布 | 矩阵热力图 |
[交易系统] --> [Kafka] --> [Flink实时计算] --> [Redis/Elasticsearch]
│
└──> [Grafana实时可视化] 组件 | 用途 | 推荐方案 |
|---|---|---|
数据采集 | 交易事件捕获 | Kafka/WebSocket |
实时计算 | 窗口聚合/地理编码 | Flink/Spark Streaming |
存储引擎 | 热力数据快速查询 | Redis(Geo+SortedSet) |
可视化 | 热力渲染 | Grafana+热力图插件 |
# 使用Docker Compose快速部署
version: '3'
services:
kafka:
image: bitnami/kafka:3.4
ports:
- "9092:9092"
flink:
image: flink:1.17
ports:
- "8081:8081"
redis:
image: redis/redis-stack:latest
ports:
- "6379:6379"
- "8001:8001" # RedisInsight
grafana:
image: grafana/grafana-enterprise:10.1
ports:
- "3000:3000" # 安装Heatmap面板插件
docker exec grafana grafana-cli \
plugins install petrslavotinek-heatmap-panel
# 安装Geomap插件(地理热力)
docker exec grafana grafana-cli \
plugins install grafana-worldmap-panel {
"trade_id": "T202311011200001",
"amount": 15000.0,
"currency": "CNY",
"geo": {
"lat": 31.2304,
"lon": 121.4737
},
"timestamp": "2023-11-01T12:00:00Z",
"product_category": "electronics"
} // 创建执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(5000);
// Kafka数据源
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers("kafka:9092")
.setTopics("trades")
.setGroupId("heatmap-processor")
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
DataStream<TradeEvent> trades = env.fromSource(
source,
WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5)),
"Kafka Source");
// 窗口聚合(每分钟/每平方公里)
trades
.assignTimestampsAndWatermarks(...)
.keyBy(event -> GeoHash.encode(event.geo.lat, event.geo.lon, 6)) // 6级Geohash
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new HeatmapAggregator())
.addSink(new RedisSink());
// 聚合函数实现
class HeatmapAggregator implements AggregateFunction<TradeEvent, HeatmapAccumulator, HeatmapData> {
public HeatmapAccumulator createAccumulator() {
return new HeatmapAccumulator();
}
public HeatmapAccumulator add(TradeEvent event, HeatmapAccumulator acc) {
acc.count++;
acc.totalAmount += event.amount;
return acc;
}
public HeatmapData getResult(HeatmapAccumulator acc) {
return new HeatmapData(acc.geohash, acc.count, acc.totalAmount);
}
} # 地理索引
GEOADD trades:geo 121.4737 31.2304 T202311011200001
# 时间窗口统计(Sorted Set)
ZADD trades:time:2023110112 1667296800 T202311011200001
# 热力值存储(Hash)
HSET trades:heatmap:2023110112:6 "wx4er5" '{"count":15,"amount":45000}' # 自动清理24小时前的数据
EXPIRE trades:heatmap:2023110112 86400 步骤:
-- 获取当前时间窗口的Geohash分布
EVAL "
local results = {}
local geohashes = redis.call('HKEYS', KEYS[1])
for _, geohash in ipairs(geohashes) do
local val = redis.call('HGET', KEYS[1], geohash)
local data = cjson.decode(val)
local lat, lon = unpack(redis.call('GEOPOS', 'trades:geo', geohash))
results[#results+1] = {
geohash = geohash,
lat = lat,
lon = lon,
count = data.count,
amount = data.amount
}
end
return results
" 1 "trades:heatmap:2023110112:6" 样式设置:
{
"colorMode": "opacity",
"heatmapRadius": 30,
"heatmapBlur": 15,
"thresholds": {
"mode": "absolute",
"steps": [
{"color": "rgba(33,102,172,0)", "value": null},
{"color": "rgb(103,169,207)", "value": 10},
{"color": "rgb(209,229,240)", "value": 50}
]
}
} 使用Heatmap面板:
SELECT
FLOOR(UNIX_TIMESTAMP(time)/60 AS time_bucket,
product_category,
COUNT(*) AS trade_count
FROM trades
WHERE $__timeFilter(time)
GROUP BY 1, 2 参数调优:
axes:
xAxis:
mode: "time"
show: true
yAxis:
decimals: 0
color:
mode: "spectrum"
fill: "dark" // 前端订阅数据更新
const subscription = {
channel: "trades/updates",
data: {
interval: "1s"
}
};
this.grafanaLive.subscribe(subscription).pipe(
map(data => this.transformToHeatmapData(data))
).subscribe(update => {
this.panelData = update;
}); func (s *TradeService) StreamTrades(ctx context.Context) {
pub := s.grafanaLive.GetPublisher()
ticker := time.NewTicker(1 * time.Second)
for {
select {
case <-ticker.C:
data := s.redisClient.HGetAll("trades:heatmap:current").Val()
pub("trades/updates", data)
case <-ctx.Done():
return
}
}
} graph LR
A[客户端内存] --> B[Redis缓存] --> C[Flink状态后端] // 在Flink中实现分层采样
.sample(
SampleStrategy.of(
rate -> Math.log(rate) > threshold, // 高频区域全采样
rate -> rate > 0.1 // 低频区域概率采样
)
) -- 使用Redis Lua脚本合并请求
local results = {}
for _, key in ipairs(KEYS) do
results[#results+1] = redis.call('HGETALL', key)
end
return results 现象 | 排查步骤 | 工具命令 |
|---|---|---|
热力点显示偏移 | 1. 检查坐标参考系(WGS84) 2. 验证Geohash精度 | GEOPOS trades:geo wx4er5 |
数据更新延迟 | 1. 检查Kafka消费延迟 2. 查看Flink Checkpoint日志 | kafka-consumer-groups --describe |
颜色映射异常 | 1. 验证阈值配置 2. 检查数据值域分布 | HGET trades:heatmap:current wx4er5 |
from sklearn.cluster import DBSCAN
# 实时聚类分析
coordinates = [[lat1,lon1], [lat2,lon2], ...]
clustering = DBSCAN(eps=0.5, min_samples=10).fit(coordinates)
# 在Grafana标注异常集群
annotations = {
"text": "异常交易集群",
"coordinates": clustering.core_sample_indices_
} // 使用Three.js插件
const heatmapTexture = new THREE.DataTexture(
heatmapData,
width,
height,
THREE.RedFormat
);
const material = new THREE.MeshBasicMaterial({map: heatmapTexture}); 项目交付物清单:
通过本方案,某证券交易所成功实现毫秒级延迟的交易热力监控,数据处理能力达到50万事件/秒。建议生产环境采用分区部署架构,在不同地域部署边缘计算节点进行预处理。
本篇的分享就到这里了,感谢观看,如果对你有帮助,别忘了点赞+收藏+关注。