Presto/Trino SQL 查询多数据源

FreeGuideOnline 最新 2026-07-09

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/query API 查看 join statistics)。
  • 尽量使用子查询预先过滤,减少参与 JOIN 的行数。
  • 如果数据量极大,考虑中间表或 ETL 预先整合。

类型映射

不同数据源的数据类型可能不完全一致。Trino 会进行自动类型映射,但偶尔会遇到类型不匹配错误。例如,MySQL 的 TINYINT(1) 默认映射为 boolean,如果期望是整数,可以在连接器配置中设置 jdbc-types-mapped-to-boolean 排除该字段。

事务与一致性

Trino 查询默认是读已提交的隔离级别,不保证跨源的数据一致性。如果业务需求严格一致性,需要在应用层处理或选择支持事务的存储系统。

实战练习:构建跨源订单分析

假设你有一个电商系统,数据分散在:

  1. MySQL(订单详情)
  2. Hive(用户画像日志)
  3. 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;