Presto/Trino SQL 查询多数据源
mermaid flowchart LR A[客户端 SQL 查询] --> B[Trino Coordinator] B --> C{Meta Store / 连接器} C --> D[(Hive)] C --> E[(MySQL)] C --> F[(Kafka)] D --> B E --> B F --> B B --> G[返回合并结果]
Coordinator 负责解析 SQL、生成执行计划,并将子任务分发给 Worker 节点。Worker 通过各自的连接器并行读取数据,最终汇总返回。
## 环境准备与架构说明
要体验多数据源查询,你需要一个运行中的 Trino 集群(至少包含一个 Coordinator 和一个 Worker)。初学者可以通过 Docker 快速搭建:
```bash
docker run -d --name trino trinodb/trino
接下来的配置均在 Coordinator 和所有 Worker 节点的 etc/catalog/ 目录下进行。
配置不同的数据源
每个数据源对应一个目录文件(etc/catalog/<catalog_name>.properties),文件内容定义了连接器的类型和连接参数。配置完成后,Trino 会自动加载,无需重启(部分版本需动态加载特性)。
配置 MySQL 数据源
创建 etc/catalog/mysql_source.properties:
connector.name=mysql
connection-url=jdbc:mysql://192.168.1.100:3306/mydb
connection-user=trino
connection-password=trino_pass
此时 Trino 中就会出现一个名为 mysql_source 的目录(catalog),mydb 内的表可通过 mysql_source.mydb.<table> 访问。
配置 Hive 数据源
Hive 连接器用于查询 HDFS 或对象存储上的数据。创建 etc/catalog/hive_source.properties:
connector.name=hive-hadoop2
hive.metastore.uri=thrift://hive-metastore:9083
如果使用本地存储,可配置 hive.config.resources 指向 core-site.xml 和 hdfs-site.xml。
配置 Kafka 数据源
Kafka 连接器将 Topic 映射为表,常用于实时流数据查询。创建 etc/catalog/kafka_source.properties:
connector.name=kafka
kafka.table-names=my_topic1,my_topic2
kafka.nodes=broker1:9092,broker2:9092
kafka.table-description-dir=etc/kafka
需要在 etc/kafka 目录下为每个表定义 JSON schema 文件。
其他常用连接器:
- PostgreSQL:
connector.name=postgresql - MongoDB:
connector.name=mongodb - Elasticsearch:
connector.name=elasticsearch - Redis:
connector.name=redis
每个连接器的详细参数请参考 Trino 官方文档。
编写跨数据源查询
配置完成后,你就可以像操作普通表一样,在一条 SQL 中混合使用不同 catalog 的表。
查询单个数据源
先验证连接是否通畅:
-- 查询 MySQL 中的订单表
SELECT * FROM mysql_source.mydb.orders LIMIT 10;
-- 查询 Hive 中的用户行为日志
SELECT * FROM hive_source.default.user_events WHERE dt = '2025-03-20' LIMIT 10;
语法:
<catalog>.<schema>.<table>。如果数据源没有 schema 概念(如 MySQL 中 schema 即 database),则 schema 就是数据库名。
跨数据源关联查询(INNER JOIN)
将 MySQL 的订单与 Hive 的用户信息关联:
SELECT
o.order_id,
o.total_amount,
u.name,
u.email
FROM
mysql_source.mydb.orders o
JOIN
hive_source.default.users u
ON o.user_id = u.user_id
WHERE
o.order_date = DATE '2025-03-20'
LIMIT 100;
Trino 会分别从 MySQL 和 Hive 拉取必要数据,在 Worker 内存中进行 JOIN 运算。
联合聚合查询(UNION ALL)
将不同数据源的统计数据合并展示:
SELECT 'web' AS source, COUNT(*) AS cnt
FROM mysql_source.mydb.page_views
WHERE date = CURRENT_DATE
UNION ALL
SELECT 'app' AS source, COUNT(*) AS cnt
FROM hive_source.default.app_events
WHERE event_date = CURRENT_DATE;
使用子查询跨源过滤
利用一个数据源的结果作为另一个数据源的过滤条件:
SELECT * FROM kafka_source.default.realtime_clicks
WHERE user_id IN (
SELECT user_id FROM mysql_source.mydb.vip_users
);
查询优化与注意事项
联邦查询虽然便捷,但不当使用可能导致性能问题。以下是一些关键要点。
谓词下推
Trino 会尽量将过滤条件下推到数据源,减少传输数据量。要确保连接器支持谓词下推(大多数关系型连接器都支持)。可以通过 EXPLAIN 查看执行计划:
EXPLAIN SELECT * FROM mysql_source.mydb.orders WHERE order_id = 123;
如果输出显示 _predicates_ 被推送到 mysql_source,说明优化生效。
谨慎使用跨源 JOIN
跨源 JOIN 的数据必须在 Trino Worker 内存中汇聚。如果两个表都非常大,会导致网络传输量和内存占用剧增。建议:
- 用小表作为驱动表(可通过
/v1/queryAPI 查看 join statistics)。 - 尽量使用子查询预先过滤,减少参与 JOIN 的行数。
- 如果数据量极大,考虑中间表或 ETL 预先整合。
类型映射
不同数据源的数据类型可能不完全一致。Trino 会进行自动类型映射,但偶尔会遇到类型不匹配错误。例如,MySQL 的 TINYINT(1) 默认映射为 boolean,如果期望是整数,可以在连接器配置中设置 jdbc-types-mapped-to-boolean 排除该字段。
事务与一致性
Trino 查询默认是读已提交的隔离级别,不保证跨源的数据一致性。如果业务需求严格一致性,需要在应用层处理或选择支持事务的存储系统。
实战练习:构建跨源订单分析
假设你有一个电商系统,数据分散在:
- MySQL(订单详情)
- Hive(用户画像日志)
- Kafka(实时点击流)
目标:查询出今日下单的高价值用户及其最近一次点击行为。
WITH
today_orders AS (
SELECT user_id, SUM(amount) AS total
FROM mysql_source.mydb.orders
WHERE order_date = CURRENT_DATE
GROUP BY user_id
HAVING SUM(amount) > 1000
),
latest_clicks AS (
SELECT
user_id,
page_url,
ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY click_time DESC) AS rn
FROM kafka_source.default.user_clicks
WHERE click_time >= CURRENT_TIMESTAMP - INTERVAL '1' HOUR
)
SELECT
o.user_id,
o.total,
u.name,
u.email,
u.region,
lc.page_url AS last_click_page
FROM today_orders o
JOIN hive_source.default.user_profiles u ON o.user_id = u.user_id
LEFT JOIN latest_clicks lc ON o.user_id = lc.user_id AND lc.rn = 1
ORDER BY o.total DESC;