本文目录导读:

很高兴为您介绍Presto的实战案例,Presto(现更名为Trino)是一个分布式SQL查询引擎,擅长对大规模数据进行快速交互式分析,下面我将从几个典型的业务场景出发,提供具体的案例和SQL示例。
电商实时大屏分析
场景描述
电商平台需要实时分析用户行为、订单情况、商品热度等指标,数据存储在Hive和Kafka中。
案例SQL
-- 1. 实时销售统计(每5分钟)
SELECT
date_trunc('minute', order_time) -
INTERVAL '5' MINUTE * (date_diff('minute', TIMESTAMP '1970-01-01', order_time) % 5) AS time_bucket,
COUNT(DISTINCT user_id) AS unique_users,
SUM(order_amount) AS total_revenue,
COUNT(*) AS total_orders,
AVG(order_amount) AS avg_order_amount
FROM hive.orders_db.orders
WHERE order_time >= current_timestamp - INTERVAL '1' HOUR
GROUP BY 1
ORDER BY 1 DESC;
-- 2. 热门商品TopN(实时)
SELECT
product_name,
category,
COUNT(*) AS sales_count,
SUM(order_amount) AS total_amount,
RANK() OVER (ORDER BY COUNT(*) DESC) AS sales_rank
FROM hive.orders_db.orders o
JOIN hive.products_db.products p ON o.product_id = p.product_id
WHERE o.order_time >= current_timestamp - INTERVAL '30' MINUTE
GROUP BY product_name, category
ORDER BY sales_count DESC
LIMIT 20;
跨数据源联邦查询
场景描述
某金融公司数据分散在MySQL(交易数据)、PostgreSQL(客户信息)、Hive(历史数据)中,需要统一查询。
案例SQL
-- 1. 客户交易行为分析(跨数据源JOIN)
SELECT
c.customer_id,
c.name,
c.level,
t.bank_card,
t.transaction_amount,
t.transaction_time,
h.previous_risk_rating
FROM mysql.finance_db.customers c
JOIN postgresql.trade_db.transactions t
ON c.customer_id = t.customer_id
JOIN hive.history_db.risk_assessment h
ON c.customer_id = h.customer_id
WHERE c.level IN ('VIP', 'SVIP')
AND t.transaction_amount > 100000
AND t.transaction_time >= current_date - INTERVAL '7' DAY
ORDER BY t.transaction_amount DESC;
-- 2. 跨数据源聚合报表
SELECT
DATE_FORMAT(t.transaction_time, '%Y-%m') AS month,
c.city,
COUNT(DISTINCT c.customer_id) AS active_customers,
SUM(t.transaction_amount) AS monthly_trading_volume,
AVG(t.transaction_amount) AS avg_trade_amount
FROM mysql.finance_db.customers c
JOIN postgresql.trade_db.transactions t ON c.customer_id = t.customer_id
GROUP BY 1, 2
HAVING monthly_trading_volume > 10000000
ORDER BY 1, 3 DESC;
日志分析平台
场景描述
某互联网公司每天产生TB级应用日志,存储在Hive中,需要快速进行异常检测、用户行为分析。
案例SQL
-- 1. 接口调用异常检测(HTTP 5xx错误分析)
SELECT
request_path AS api_endpoint,
COUNT(*) AS total_requests,
SUM(CASE WHEN status_code >= 500 THEN 1 ELSE 0 END) AS error_count,
ROUND(SUM(CASE WHEN status_code >= 500 THEN 1 ELSE 0 END) * 100.0 / COUNT(*), 2) AS error_rate,
AVG(response_time_ms) AS avg_response_time,
PERCENTILE_ESTIMATE(response_time_ms, 0.95) AS p95_response_time
FROM hive.logs_db.access_logs
WHERE log_date >= current_date - INTERVAL '1' DAY
GROUP BY request_path
HAVING error_rate > 1.0
ORDER BY error_rate DESC;
-- 2. 用户访问路径分析(会话化分析)
WITH session_data AS (
SELECT
user_id,
request_path,
request_time,
DATE_DIFF('MINUTE',
LAG(request_time) OVER (PARTITION BY user_id ORDER BY request_time),
request_time
) AS time_gap
FROM hive.logs_db.access_logs
WHERE log_date = current_date - INTERVAL '1' DAY
),
session_boundaries AS (
SELECT *,
SUM(CASE WHEN time_gap > 30 OR time_gap IS NULL THEN 1 ELSE 0 END)
OVER (PARTITION BY user_id ORDER BY request_time) AS session_id
FROM session_data
)
SELECT
user_id,
session_id,
MIN(request_time) AS session_start,
MAX(request_time) AS session_end,
DATEDIFF('MINUTE', MIN(request_time), MAX(request_time)) AS session_duration_min,
ARRAY_AGG(request_path ORDER BY request_time) AS page_sequence
FROM session_boundaries
GROUP BY user_id, session_id
HAVING DATEDIFF('MINUTE', MIN(request_time), MAX(request_time)) > 5
ORDER BY session_start DESC LIMIT 1000;
广告投放效果分析
场景描述
广告公司需要分析不同渠道、不同创意形式的广告投放效果,数据存储在Hive中。
案例SQL
-- 1. 广告效果多维度分析
SELECT
channel,
creative_type,
campaign_name,
COUNT(DISTINCT imp_id) AS impressions,
COUNT(DISTINCT click_id) AS clicks,
COUNT(DISTINCT conv_id) AS conversions,
ROUND(COUNT(DISTINCT click_id) * 100.0 / COUNT(DISTINCT imp_id), 2) AS ctr,
ROUND(COUNT(DISTINCT conv_id) * 100.0 / COUNT(DISTINCT click_id), 2) AS cvr,
ROUND(SUM(cost) / COUNT(DISTINCT conv_id), 2) AS cac,
ROUND(SUM(revenue) - SUM(cost), 2) AS profit
FROM hive.advertising_db.ad_performance
WHERE impression_date >= current_date - INTERVAL '7' DAY
GROUP BY 1, 2, 3
HAVING impressions > 10000
ORDER BY profit DESC;
-- 2. 渠道归因分析(最后一次点击归因)
WITH last_click_attribution AS (
SELECT
user_id,
conv_id,
LAST_VALUE(channel) OVER (
PARTITION BY user_id, conv_id
ORDER BY click_time
ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
) AS attributed_channel,
conversion_time,
revenue
FROM hive.advertising_db.conversions c
JOIN hive.advertising_db.clicks cl ON c.user_id = cl.user_id
AND cl.click_time <= c.conversion_time
AND cl.click_time >= c.conversion_time - INTERVAL '30' DAY
)
SELECT
attributed_channel,
COUNT(DISTINCT conv_id) AS attributed_conversions,
SUM(revenue) AS attributed_revenue,
RANK() OVER (ORDER BY SUM(revenue) DESC) AS channel_rank
FROM last_click_attribution
WHERE conversion_time >= current_date - INTERVAL '30' DAY
GROUP BY 1
ORDER BY 2 DESC;
性能优化案例
场景描述
面对大型表查询出现慢查询,需要进行针对性优化。
-- 1. 使用WITH子句优化复杂查询
WITH filtered_orders AS (
-- 先过滤,减少后续JOIN数据量
SELECT order_id, user_id, amount, order_time
FROM hive.orders_db.orders
WHERE order_date = current_date - INTERVAL '1' DAY
),
user_segments AS (
SELECT user_id, segment
FROM hive.users_db.users
WHERE segment IN ('A', 'B')
)
SELECT
us.segment,
DATE_TRUNC('HOUR', fo.order_time) AS hour,
COUNT(*) AS order_count,
SUM(fo.amount) AS total_amount
FROM filtered_orders fo
JOIN user_segments us ON fo.user_id = us.user_id
GROUP BY 1, 2;
-- 2. 使用分区字段过滤
-- 错误:没有利用分区
SELECT * FROM hive.logs_db.access_logs
WHERE log_date >= date '2024-01-01' AND log_date < date '2024-01-02';
-- 正确:利用分区字段,查询性能提升10-100倍
SELECT * FROM hive.logs_db.access_logs
WHERE log_date >= date '2024-01-01' AND log_date < date '2024-01-02';
-- 3. 使用物化视图(需要安装Materialized View插件)
CREATE MATERIALIZED VIEW hive.data_mart.daily_sales_mv AS
SELECT
order_date,
product_category,
COUNT(DISTINCT user_id) AS unique_buyers,
SUM(order_amount) AS total_sales,
COUNT(*) AS order_count
FROM hive.orders_db.orders o
JOIN hive.products_db.products p ON o.product_id = p.product_id
GROUP BY 1, 2;
-- 查询物化视图(自动使用最新预计算结果)
SELECT * FROM hive.data_mart.daily_sales_mv
WHERE order_date >= current_date - INTERVAL '3' MONTH;
连接器配置
catalogs/
├── mysql.properties # MySQL连接器配置
├── hive.properties # Hive连接器配置
├── kafka.properties # Kafka连接器配置
└── postgresql.properties # PostgreSQL连接器配置
优化建议
- 分区裁剪:始终使用分区字段过滤
- 谓词下推:WHERE条件尽量写在JOIN之前
- 数据本地性:协调计算和存储节点尽量在同一网络
- 适当调整并发:根据集群资源调整task和node数量
监控指标
查询执行时间 | 内存使用 | CPU使用率 | 数据扫描量 | 网络传输量
这些案例展示了Presto在跨数据源分析、实时查询、日志处理等场景的强大能力,如果您有具体的业务场景,欢迎进一步讨论!