Towards Data Science

Running SQL Concurrently Across Three Remote DuckDB Servers with Quack

8.5内容质量

TL;DR · AI 摘要

Quack协议实现跨三台远程DuckDB服务器的SQL并发执行,通过AWS EC2实验验证可行性。

核心要点

  • Quack协议通过HTTP实现远程DuckDB数据库通信,支持跨服务器读写操作。
  • 实验部署3台AWS EC2服务器,使用ARM64架构和DuckDB 1.5.5版本。
  • 并发SQL执行需配置Quack扩展、systemd服务及自动关机计时器。

结构提纲

按章节快速跳转。

  1. 介绍Quack协议的远程数据库通信能力及实验目标。

  2. 通过CloudFormation部署3台AWS EC2服务器并安装依赖组件。

  3. 详细说明Quack扩展、服务配置及认证令牌设置方法。

  4. 展示如何通过协调节点并行执行跨服务器SQL查询。

  5. 对比单节点与多节点执行性能差异及潜在优化方向。

思维导图

用一张图看清主题之间的关系。

查看大纲文本(无障碍 / 无 JS 友好)
  • Quack远程SQL并发实验
    • 核心组件
      • DuckDB 1.5.5
      • Quack协议
      • AWS EC2
    • 部署流程
      • CloudFormation配置
      • Python环境搭建
      • systemd服务设置
    • 实验验证
      • 并发SQL执行
      • 性能对比
      • 自动关机机制

金句 / Highlights

值得收藏与分享的关键句。

#DuckDB#Quack#SQL并发#AWS EC2#数据工程
打开原文

通过 Quack 在三个远程 DuckDB 服务器上并发运行 SQL | Towards Data Science

数据工程

通过 Quack 在三个远程 DuckDB 服务器上并发运行 SQL

一项关于远程 SQL 执行的小型实验

Thomas Reid

2026 年 8 月 16 日

16 分钟阅读时间

分享

AI 生成图片

几个月前,DuckDB 的开发团队发布了一个名为 Quack 的数据库通信协议。该协议的主要目标是让位于不同服务器上的 DuckDB 数据库可以通过 HTTP 以客户端/服务器架构进行通信,并实现彼此数据的读写操作。

换句话说,通过 Quack 协议,位于服务器 A 上的 DuckDB 现在可以查询或写入远程服务器 B 上的 DuckDB 数据库。这听起来可能有点像分布式数据处理,但实际上并非如此。DuckDB 团队特别强调 Quack 并不支持分布式查询处理。

尽管如此,我仍然对此感到好奇,并看到了 Quack 的许多潜在用途。特别是,我很好奇是否可以向每个服务器同时发出并行的 SQL 语句,并让每个查询的结果在协调服务器上被收集和展示。

请注意,这与在单个 SQL 语句中跨不同服务器连接表是不同的概念。例如,DuckDB 可以通过将服务器 B 上的数据库附加到服务器 A,然后在服务器 A 上运行 SQL 语句来读取服务器 B 上的数据,就像这些数据是本地数据一样。

为了研究 Quack 应该支持的并发读写操作,我创建了一个名为 cluster-duck 的 GitHub 仓库。需要说明的是,这不是一个分布式 DuckDB 集群。这只是对您可能已经知道的常见粗俗表达的一种幽默诠释。但它确实可以让您在远程 DuckDB 数据库上执行并发读写操作。

注意:在继续之前,我想说明一下,我与本文中提到的任何产品、系统或其创建者都没有任何关联或商业关系。

环境搭建

为了测试 Quack,我通过 CloudFormation 堆栈部署了 3 台 AWS EC2 服务器。每台服务器都运行一个 DuckDB 数据库,其中一台服务器还充当协调节点。您可以在我的 GitHub 仓库中找到 CF 堆栈。

从 AWS CloudShell(或如果您安装了 AWS CLI,也可以在本地)部署堆栈,命令如下:

code
aws cloudformation deploy \
--region us-east-2 \
--stack-name cluster-duck-test-v2 \
--template-file python-reference/infra/aws/cluster-duck-3-node.yaml \
--capabilities CAPABILITY_NAMED_IAM

对于所有三台 EC2 服务器,安装了以下组件:

  • Amazon Linux 2023 ARM64
  • 512 MB 交换文件
  • Python 3.12 和 pip
  • Python 虚拟环境位于:/opt/cluster-duck-venv
  • duckdb==1.5.5
  • boto3
  • DuckDB 1.5.5 ARM64 CLI 位于:/usr/local/bin/duckdb
  • 官方 DuckDB Quack 扩展,由工作进程加载
  • Quack 服务器程序位于:/opt/cluster-duck/quack_server.py
  • 工作数据库位于:/var/lib/cluster-duck/worker.duckdb
  • systemd 服务名为:cluster-duck-quack.service
  • 自动关机计时器,默认为 4 小时

Quack 默认监听 9494 端口。其认证令牌从加密的 SSM 参数存储参数中获取。

Python 源代码被复制到所有三台服务器上,因为它们共享相同的 CloudFormation 启动模板。但是,协调器仅在一台服务器上被设置为可执行命令——通常是工作节点 1。

这些文件安装在每台服务器上:

code
/opt/cluster-duck/quack_server.py
/opt/cluster-duck/seed_related_data.py
/opt/cluster-duck/related_cluster_sql.py
/opt/cluster-duck/cluster_duck/sql_api.py

协调服务器(Worker 1)还会获得这些命令启动器:

code
/usr/local/bin/cluster-duck
/usr/local/bin/cluster-duck-sql

在协调器上运行 cluster-duck-sql 时,最终会启动:

code
/opt/cluster-duck-venv/bin/python
/opt/cluster-duck/related_cluster_sql.py

Python 源代码被压缩并直接嵌入到 CloudFormation 模板中,作为 Base64 编码的存档。

在 EC2 引导过程中,用户数据脚本会执行以下操作:

  • 解码嵌入的存档。
  • 创建 /opt/cluster-duck 目录。
  • 将 Python 文件提取到该目录中。
  • 在 worker 1 上创建命令启动器。
  • 在每个 worker 上启动 Quack 服务。

所有需要的内容都包含在 CloudFormation 文件中。

协调机制的工作原理

Quack 会将每个 SQL 语句发送到选定的 DuckDB 服务器,并返回结果。协调三个调用的部分是运行在 Worker 1 上的 Python 代码。

首先,每个 —query 或 —query-file 参数都会被验证为单个 SQL 语句,并转换为 QueryFragment。这些片段会按照提供的顺序进行标记:

code
fragments.append(
QueryFragment(worker_id, f"query-{index}", sql)
)

然后协调器为每个片段创建一个线程,并创建一个参与人数相同的屏障:

code
barrier = threading.Barrier(len(fragments))
epoch = time.perf_counter()

def run_fragment(fragment):
    barrier.wait()
    started = time.perf_counter()
    raw_result = self.executor(fragment.worker_id, fragment.sql)
    finished = time.perf_counter()
    return {
        "start_offset_ms": (started - epoch) * 1000,
        "duration_ms": (finished - started) * 1000,
        "result": raw_result,
    }
with ThreadPoolExecutor(max_workers=len(fragments)) as pool:
    futures = {
        pool.submit(run_fragment, fragment): fragment
        for fragment in fragments
    }

屏障会一直保持线程,直到所有片段都准备就绪,然后同时释放它们。由于操作系统调度仍然适用,这些线程不会在完全相同的 CPU 周期启动,这就是为什么输出包含测量的启动时间差。

在协调器上,每个片段都会获得自己的本地 DuckDB 客户端连接。该连接加载 Quack,附加一个远程 worker,并通过 remote.query() 发送语句:

code
with duckdb.connect() as connection:
    connection.execute("INSTALL quack")
    connection.execute("LOAD quack")
    connection.execute(
        f"ATTACH {endpoint} AS remote "
        f"(TYPE quack, TOKEN {token}, DISABLE_SSL true)"
    )
    cursor = connection.execute(
        f"SELECT * FROM remote.query({sql_string(sql)})"
     )
    columns = tuple(description[0] for description in cursor.description)
    return FragmentResult(columns, cursor.fetchall())

调用 fetchall() 会在 Worker 1 上生成每个结果。协调器等待所有 future 完成,记录每个 future 的启动时间和耗时,然后将各个结果合并输出。Quack 负责远程执行和传输;片段构建、同时释放、计时和结果收集均由 Python 完成。

创建我们的测试数据

每个服务器都有不同的 DuckDB 数据库配置如下。

code
Worker      数据库文件                        生成的表
---------------------------------------------------------------
Worker 1   /var/lib/cluster-duck/worker.duckdb   sales
Worker 2   /var/lib/cluster-duck/worker.duckdb   customers
Worker 3   /var/lib/cluster-duck/worker.duckdb   products

每个生成的表都包含适当的人工生成数据集,包含一千万条记录。以下是每个表的前5条记录,帮助您更好地了解表中内容。

code
query-1 (worker-1) - select * from sales limit 5

sale_id  customer_id  product_id  quantity  sales_channel  payment_method  sale_status  catalogue_price  sold_unit_price  discount_pct  sale_date
-------  -----------  ----------  --------  -------------  --------------  -----------  ---------------  ---------------  ------------  ----------
1        7,920        104,730     2         store          bank_transfer   shipped      1,052.3          999.69           0.05          2024-01-02
2        15,839       209,459     3         marketplace    wallet          processing   99.59            89.63            0.1           2024-01-03
3        23,758       314,188     4         telephone      invoice         returned     1,146.88         974.85           0.15          2024-01-04
4        31,677       418,917     5         online         card            cancelled    194.17           194.17           0             2024-01-05
5        39,596       523,646     1         store          bank_transfer   completed    1,241.46         1,179.39         0.05          2024-01-06

query-2 (worker-2) - select * from customers limit 5

customer_id  customer_code    country  segment         membership_tier  is_active  credit_limit  joined_date  last_seen_at
-----------  ---------------  -------  --------------  ---------------  ---------  ------------  -----------  -------------------
1            CUST-0000000001  US       small_business  silver           True       250.10        2015-01-02   2025-01-01 00:00:01
2            CUST-0000000002  DE       enterprise      gold             True       250.20        2015-01-03   2025-01-01 00:00:02
3            CUST-0000000003  FR       public_sector   standard         True       250.30        2015-01-04   2025-01-01 00:00:03
4            CUST-0000000004  CA       consumer        silver           True       250.40        2015-01-05   2025-01-01 00:00:04
5            CUST-0000000005  AU       small_business  gold             True       250.50        2015-01-06   2025-01-01 00:00:05

query-3 (worker-3) - select * from products limit 5

product_id  sku             category  brand    supplier_region  catalogue_price  stock_quantity  discontinued  introduced_date
----------  --------------  --------  -------  ---------------  ---------------  --------------  ------------  ---------------
1           SKU-0000000001  home      Bramble  EU               5.01             13              False         2020-01-02
2           SKU-0000000002  garden    Cobalt   US               5.02             26              False         2020-01-03
3           SKU-0000000003  sports    Dove     APAC             5.03             39              False         2020-01-04
4           SKU-0000000004  clothing  Elm      UK               5.04             52              False         2020-01-05
5           SKU-0000000005  food      Aster    EU               5.05             65              False         2020-01-06

数据是在第一次 EC2 引导期间直接在每个 DuckDB 数据库内部生成的。它不会从您的计算机上传,也不会在服务器之间复制。每台服务器都遵循以下流程。

1. 确定当前是哪台工作节点

CloudFormation 为每个 EC2 实例分配了一个 WorkerIndex 标签:

code
Worker 1 → WorkerIndex=1
Worker 2 → WorkerIndex=2
Worker 3 → WorkerIndex=3
引导脚本通过 EC2 实例元数据服务读取该标签:
WORKER_INDEX=$(curl -fsS \
  -H "X-aws-ec2-metadata-token: $IMDS_TOKEN" \
  http://169.254.169.254/latest/meta-data/tags/instance/WorkerIndex)

2. 运行数据生成程序

CloudFormation 在每台服务器上安装了这个程序:

code
/opt/cluster-duck/seed_related_data.py

然后执行:

code
/opt/cluster-duck-venv/bin/python /opt/cluster-duck/seed_related_data.py \
--worker "$WORKER_INDEX" \
--rows 10000000

行数来自 CloudFormation 的 RowCount 参数,该参数默认值为 1000 万。

3. 打开工作节点的 DuckDB 文件

该程序打开:

code
/var/lib/cluster-duck/worker.duckdb

其中包含以下代码:

code
with duckdb.connect(str(args.database)) as connection:
    connection.execute(
        CREATE_SQL[args.worker],
        {"row_count": args.rows},
    )

每台服务器使用相同的数据库文件名,但不同 EC2 实例上的文件是不同的。

访问协调器终端窗口

我们希望运行一些演示,为此您需要能够访问协调器 EC2 服务器的 CLI 终端。为此,请打开 AWS 控制台并进入 EC2 控制台。您将看到如下界面。

点击对应您协调器实例的实例 ID。在下一页的右上角会有连接按钮。点击该按钮,您将看到如下界面。

确保您已选择 SSM 会话管理器单选按钮,然后点击屏幕右下角的连接按钮。这样应该可以访问 CLI 终端窗口,如下所示:

示例

在以下示例中,为了尽可能清晰,我使用了原始 SQL 文本作为代码片段内容,但也可以将 SQL 存储在单独的文件中并使用这些文件作为输入。例如:

code
[root@ip-10-42-0-10 ~]# cluster-duck-sql \
  --query-file "worker-1=/root/cluster-duck-sql/sales.sql" \
  --query-file "worker-2=/root/cluster-duck-sql/customers.sql" \
  --query-file "worker-3=/root/cluster-duck-sql/products.sql"

1. 执行一些简单的 SQL 语句

在终端 CLI 中,输入以下代码:

code
sh-5.2$ sudo -i
[root@ip-10-42-0-10 ~]# cluster-duck-sql \
  --query "worker-1=SELECT sale_status, COUNT(*) FROM sales GROUP BY sale_status ORDER BY sale_status" \
  --query "worker-2=SELECT country, COUNT(*) FROM customers GROUP BY country ORDER BY country" \
  --query "worker-3=SELECT category, COUNT(*) FROM products GROUP BY category ORDER BY category"

并发远程查询

worker 表名 开始偏移毫秒 持续时间秒 -------- ------- --------------- ---------------- worker-1 query-1 0.996 0.572 worker-2 query-2 5.514 0.572 worker-3 query-3 0.738 0.54 开始时间差: 4.776 毫秒 query-1 (worker-1) sale_status count_star() ----------- ------------ cancelled 2,000,000 completed 2,000,000 processing 2,000,000 returned 2,000,000 shipped 2,000,000 query-2 (worker-2) country count_star() ------- ------------ AU 1,666,666 CA 1,666,667 DE 1,666,667 FR 1,666,667 UK 1,666,666 US 1,666,667 query-3 (worker-3) category count_star() ----------- ------------ clothing 1,666,667 electronics 1,666,666 food 1,666,666 garden 1,666,667 home 1,666,667 sports 1,666,667

code

### 2. 一些复杂的 SQL(我删减了一些输出以节省空间)

[root@ip-10-42-0-10 ~]# time cluster-duck-sql \ --query "worker-1=WITH daily AS ( SELECT sale_date, sales_channel, payment_method, sale_status, COUNT(*) AS transaction_count, SUM(quantity) AS units, SUM(quantity * sold_unit_price) AS revenue, AVG(discount_pct) AS average_discount, QUANTILE_CONT(sold_unit_price, 0.50) AS median_price, QUANTILE_CONT(sold_unit_price, 0.95) AS p95_price FROM sales GROUP BY ALL ), analysed AS ( SELECT *, SUM(revenue) OVER ( PARTITION BY sales_channel ORDER BY sale_date ROWS BETWEEN 29 PRECEDING AND CURRENT ROW ) AS rolling_30_row_revenue, RANK() OVER ( PARTITION BY sale_date ORDER BY revenue DESC ) AS daily_revenue_rank FROM daily ) SELECT * FROM analysed WHERE daily_revenue_rank <= 3 ORDER BY sale_date DESC, daily_revenue_rank LIMIT 100" \ --query "worker-2=WITH customer_groups AS ( SELECT country, segment, membership_tier, is_active, YEAR(joined_date) AS joined_year, CASE WHEN credit_limit < 2500 THEN 'under_2500' WHEN credit_limit < 5000 THEN '2500_to_4999' WHEN credit_limit < 7500 THEN '5000_to_7499' ELSE '7500_plus' END AS credit_band, COUNT(*) AS customer_count, AVG(credit_limit) AS average_credit_limit, STDDEV_POP(credit_limit) AS credit_limit_stddev, QUANTILE_CONT(credit_limit, 0.50) AS median_credit_limit, QUANTILE_CONT(credit_limit, 0.95) AS p95_credit_limit, MIN(joined_date) AS first_joined, MAX(last_seen_at) AS most_recent_activity FROM customers GROUP BY ALL ), ranked AS ( SELECT *, SUM(customer_count) OVER ( PARTITION BY country ) AS country_total, RANK() OVER ( PARTITION BY country ORDER BY customer_count DESC ) AS group_rank FROM customer_groups ) SELECT *, ROUND(100.0 * customer_count / country_total, 2) AS percentage_of_country FROM ranked WHERE group_rank <= 10 ORDER BY country, group_rank LIMIT 100" \ --query "worker-3=WITH inventory_groups AS ( SELECT category, brand, supplier_region, discontinued, YEAR(introduced_date) AS introduced_year, COUNT(*) AS product_count, SUM(stock_quantity) AS stock_units, SUM(stock_quantity * catalogue_price) AS inventory_value, AVG(catalogue_price) AS average_price, STDDEV_POP(catalogue_price) AS price_stddev, QUANTILE_CONT(catalogue_price, 0.50) AS median_price, QUANTILE_CONT(catalogue_price, 0.95) AS p95_price FROM products GROUP BY ALL ), ranked AS ( SELECT *, SUM(inventory_value) OVER ( PARTITION BY category ) AS category_inventory_value, RANK() OVER ( PARTITION BY category ORDER BY inventory_value DESC ) AS inventory_rank FROM inventory_groups ) SELECT *, ROUND( 100.0 * inventory_value / category_inventory_value, 2 ) AS percentage_of_category_value

code

FROM ranked
  WHERE inventory_rank <= 10
  ORDER BY category, inventory_rank
  LIMIT 100"

#
# 输出
并发远程查询
工作线程    表    开始偏移时间(毫秒)  持续时间(秒)
--------  -------  ---------------  ----------------
worker-1  query-1  1.673            15.491
worker-2  query-2  1.84             4.045
worker-3  query-3  1.42             4.045
开始时间差:0.420 毫秒

query-1 (worker-1)
销售日期   销售渠道  支付方式  销售状态  交易数量  单位数   收入        平均折扣  中位价格  95百分位价格  30天滚动收入  每日收入排名
----------  -------------  --------------  -----------  -----------------  ------  -------------  ----------------  ------------  ---------  ----------------------  ------------------
2025-12-30  store          bank_transfer   cancelled    6,849              34,245  32,722,337.65  0.05              955.72        1,810.568  588,590,148.25          1
...
...
...
2025-11-11  marketplace    wallet          completed    6,849              6,849   6,199,619.04   0.1               905.5         1,714.096  557,458,798.88          2

query-2 (worker-2)
国家  客户分层         会员等级  是否激活  注册年份  信用区间  客户数量  平均信用额度  信用额度标准差  中位信用额度  95百分位信用额度  首次注册  最近活动时间  国家总数量  分组排名  占国家比例
-------  --------------  ---------------  ---------  -----------  -----------  --------------  --------------------  -------------------  -------------------  ----------------  ------------  --------------------  -------------  ----------  ---------------------
AU       small_business  gold             True       2,017        7500_plus    21,659          8,874.317             794.81               8878.10              10115.70          2017-01-01    2025-04-26 17:20:41   1,666,666      1  1.3
AU       public_sector   gold             True       2,017        7500_plus    21,649          8,873.263             794.868              8876.70              10114.30          2017-01-01    2025-04-26 17:20:35   1,666,666      2  1.3
...
...
...
US       small_business  silver           True       2,022        7500_plus    21,605          8,876.671             792.939              8873.30              10110.10          2022-01-01    2025-04-26 17:46:37   1,666,667      10  1.3

query-3 (worker-3)
category     brand    supplier_region  discontinued  introduced_year  product_count  stock_units  inventory_value  average_price  price_stddev  median_price  p95_price  category_inventory_value  inventory_rank  percentage_of_category_value
-----------  -------  ---------------  ------------  ---------------  -------------  -----------  ---------------  -------------  ------------  ------------  ---------  ------------------------  --------------  ----------------------------
服装         Aster    US               False         2,020            33,457         83,743,010   84191813103.00   1,004.833      577.373       1004.90       1904.90    4188462626332.28          1               2.01
服装         Aster    UK               False         2,020            33,458         83,228,180   83715767600.00   1,005.04       577.369       1005.20       1905.42    4188462626332.28          2               2
...
...
...
家居         Dove     APAC             False         2,024            32,999         82,797,401   83251410292.23   1,004.931      577.217       1005.23       1904.83    4190160306717.15          8               1.99
家居         Dove     APAC             False         2,020            33,008         82,760,072   83227553226.96   1,005.078      577.503       1005.33       1905.43    4190160306717.15          9               1.99
...
...
运动         Elm      EU               False         2,020            33,006         82,731,982   83205651717.58   1,005.109      577.373       1005.19       1905.44    4190156973777.15          9               1.99
运动         Bramble  EU               False         2,020            33,005         82,700,285   83167881906.05   1,004.964      577.415       1005.21       1905.36    4190156973777.15          10              1.98
real    0m17.664s
user    0m1.643s
sys     0m0.389s
[root@ip-10-42-0-10 ~]#

#### 3. 并发写入/读取

为展示这一点,我们将同时向销售表中写入20条新记录并发起三次读取操作。每次读取都能看到一个一致性快照,其中包含在该读取开始前已提交的独立插入事务。因此,它可能看不到新行、看到部分新行或看到所有新行,但绝不会看到某个插入操作的中途状态或在SELECT扫描过程中发生变化的结果。

本例中的每个插入都是一个独立的自动提交事务。如果19个插入成功而1个失败,那么19个成功的写入操作将保持提交状态。这里不存在分布式事务,Cluster-Duck不会回滚成功的语句。

为了使演示可重复,首先删除任何早期运行留下的行:

code
[root@ip-10-42-0-10 ~]# cluster-duck-sql \
  --allow-write \
  --query "worker-1=DELETE FROM sales WHERE sale_id BETWEEN 30000001 AND 30000020"

还需要注意,为了对数据库中的数据进行修改,我们应该提供—allow-write参数。

现在一起运行20次插入和三次读取。要查看输出中每个查询标签运行的SQL语句,可以使用—show-sql参数。

code
[root@ip-10-42-0-10 ~]# cluster-duck-sql \
  --show-sql \
  --allow-write \
  --query "worker-1=INSERT INTO sales VALUES (30000001,1,1,1,'online','card','completed',100.00,100.00,0.0,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000002,2,2,2,'store','bank_transfer','processing',110.00,104.50,0.05,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000003,3,3,3,'marketplace','wallet','shipped',120.00,108.00,0.10,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000004,4,4,4,'telephone','invoice','completed',130.00,110.50,0.15,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000005,5,5,5,'online','card','processing',140.00,112.00,0.20,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000006,6,6,1,'store','bank_transfer','shipped',150.00,150.00,0.0,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000007,7,7,2,'marketplace','wallet','completed',160.00,152.00,0.05,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000008,8,8,3,'telephone','invoice','processing',170.00,153.00,0.10,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000009,9,9,4,'online','card','shipped',180.00,153.00,0.15,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000010,10,10,5,'store','bank_transfer','completed',190.00,152.00,0.20,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000011,11,11,1,'marketplace','wallet','processing',200.00,200.00,0.0,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000012,12,12,2,'telephone','invoice','shipped',210.00,199.50,0.05,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000013,13,13,3,'online','card','completed',220.00,198.00,0.10,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000014,14,14,4,'store','bank_transfer','processing',230.00,195.50,0.15,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000015,15,15,5,'marketplace','wallet','shipped',240.00,192.00,0.20,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000016,16,16,1,'telephone','invoice','completed',250.00,250.00,0.0,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000017,17,17,2,'online','card','processing',260.00,247.00,0.05,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000018,18,18,3,'store','bank_transfer','shipped',270.00,243.00,0.10,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000019,19,19,4,'marketplace','wallet','completed',280.00,238.00,0.15,DATE '2026-08-09')" \
  --query "worker-1=INSERT INTO sales VALUES (30000020,20,20,5,'telephone','invoice','processing',290.00,232.00,0.20,DATE '2026-08-09')" \
  --query "worker-1=SELECT COUNT(*) AS visible_rows FROM sales WHERE sale_id BETWEEN 30000001 AND 30000020" \
  --query "worker-1=SELECT COUNT(*) AS visible_rows FROM sales WHERE sale_id BETWEEN 30000001 AND 30000020" \
  --query "worker-1=SELECT COUNT(*) AS visible_rows FROM sales WHERE sale_id BETWEEN 30000001 AND 30000020"

#
# 输出

 /think

并发远程查询
worker    table     start_offset_ms  duration_seconds
--------  --------  ---------------  ----------------
worker-1  query-19  77.923           3.648
worker-1  query-11  26.242           3.705
worker-1  query-23  4.64             3.727
worker-1  query-5   8.747            3.691
worker-1  query-1   4.832            3.744
worker-1  query-2   5.243            3.903
worker-1  query-15  113.315          3.823
worker-1  query-6   15.126           3.923
worker-1  query-22  84.907           3.853
worker-1  query-14  35.127           3.907
worker-1  query-18  71.325           3.871
worker-1  query-21  56.68            3.888
worker-1  query-12  28.05            3.936
worker-1  query-7   15.349           3.949
worker-1  query-20  89.082           3.875
worker-1  query-10  36.901           3.93
worker-1  query-16  77.02            3.895
worker-1  query-13  19.964           3.954
worker-1  query-3   5.632            3.968
worker-1  query-8   58.938           3.916
worker-1  query-9   15.774           3.959
worker-1  query-4   50.36            3.954
worker-1  query-17  84.575           3.936

开始时间差:108.675 毫秒

query-19 (worker-1)
SQL:
INSERT INTO sales VALUES (30000019,19,19,4,'marketplace','wallet','completed',280.00,238.00,0.15,DATE '2026-08-09')

结果:
计数
-----
1

query-11 (worker-1)
SQL:
INSERT INTO sales VALUES (30000011,11,11,1,'marketplace','wallet','processing',200.00,200.00,0.0,DATE '2026-08-09')

结果:
计数
-----
1

query-23 (worker-1)
SQL:
SELECT COUNT(*) AS visible_rows FROM sales WHERE sale_id BETWEEN 30000001 AND 30000020

结果:
visible_rows
------------
4
...
...
query-22 (worker-1)
SQL:
SELECT COUNT(*) AS visible_rows FROM sales WHERE sale_id BETWEEN 30000001 AND 30000020

结果:
visible_rows
------------
16
...
...

query-18 (worker-1)
SQL:
INSERT INTO sales VALUES (30000018,18,18,3,'store','bank_transfer','shipped',270.00,243.00,0.10,DATE '2026-08-09')

结果:
计数
-----
1

query-21 (worker-1)
SQL:
SELECT COUNT(*) AS visible_rows FROM sales WHERE sale_id BETWEEN 30000001 AND 30000020

结果:
visible_rows
------------
13
...
...
query-17 (worker-1)
SQL:
INSERT INTO sales VALUES (30000017,17,17,2,'online','card','processing',260.00,247.00,0.05,DATE '2026-08-09')

结果:
计数
-----
1
[root@ip-10-42-0-10 ~]#

输出显示,当query-23执行时,已有4条记录被插入。当query-22执行时,已有16条记录被插入,而query-21执行时有13条新记录被写入。这符合逻辑,因为我们从start_offset_ms时间可以看出查询的执行顺序是query-23,然后是query-21,最后是query-22。

5. 你也可以运行DDL语句

在每个数据库中创建3个新表,然后查询它们。

code
[root@ip-10-42-0-10 ~]# cluster-duck-sql \
  --allow-write \
  --query "worker-1=CREATE OR REPLACE TABLE sales_agg AS SELECT sale_status, COUNT(*) AS sale_count FROM sales WHERE sale_status = 'cancelled' GROUP BY sale_status" \
  --query "worker-2=CREATE OR REPLACE TABLE customers_agg AS SELECT country, COUNT(*) AS customer_count FROM customers WHERE country = 'UK'  GROUP BY country" \
  --query "worker-3=CREATE OR REPLACE TABLE products_agg AS SELECT category, COUNT(*) AS product_count FROM products WHERE category = 'sports' GROUP BY category"

 /think

cluster-duck-sql \
  --query "worker-1=SELECT * FROM sales_agg ORDER BY sale_count DESC" \
  --query "worker-2=SELECT * FROM customers_agg ORDER BY customer_count DESC" \
  --query "worker-3=SELECT * FROM products_agg ORDER BY product_count DESC"

#
# 输出

并发远程查询
worker    table    start_offset_ms  duration_seconds
--------  -------  ---------------  ----------------
worker-1  query-1  1.141            0.599
worker-2  query-2  1.033            0.603
worker-3  query-3  0.755            0.622
启动时间差: 0.386 毫秒
query-1 (worker-1)
计数
-----
1
query-2 (worker-2)
计数
-----
1
query-3 (worker-3)
计数
-----
1
并发远程查询
worker    table    start_offset_ms  duration_seconds
--------  -------  ---------------  ----------------
worker-1  query-1  0.925            0.464
worker-2  query-2  1.178            0.448
worker-3  query-3  0.675            0.467
启动时间差: 0.503 毫秒
query-1 (worker-1)
sale_status  sale_count
-----------  ----------
cancelled    2,000,000
query-2 (worker-2)
country  customer_count
-------  --------------
UK       1,666,666
query-3 (worker-3)
category  product_count
--------  -------------
sports    1,666,667

所有这些的成本

这种设置的成本不应成为问题。首先,DuckDB 和 Quack 都可以免费下载和使用。我所部署的三台 EC2 服务器是微型的 t4g.nano 实例。除此之外,我们还有三个 8GB 的 gp3 卷、公共 IPv4 地址、Systems Manager、Parameter Store 和一个 Lambda 自定义资源。Lambda 调用是短暂的,但会一直部署到堆栈被删除。这个 Lambda 在 CloudFormation 模板中被命名为 TokenManagerFunction,它的唯一任务是管理三个 Quack 认证令牌。其工作流程如下:

code
创建 CloudFormation 堆栈
          ↓
调用 TokenManager Lambda
          ↓
生成三个随机的 64 字符令牌
          ↓
将它们存储为 SSM 中的 SecureString 参数
          ↓
返回成功并停止

如果我们运行这个完整设置 4 小时,预估成本如下:

code
组件                     预估成本
4 小时 EBS               $0.011
4 小时 EC2               $0.050
4 小时公共 IPv4         $0.060

总计                     $0.121

$0.121 是 us-east-2 区域的预估成本,未包含任何信用额度、免费套餐优惠、税费和数据传输费用。符合条件的 AWS 免费套餐用户可能可以免费获得部分公共 IPv4 小时。AWS 按秒计费 gp3 存储,且最低 60 秒。

尽管如此,为了安心起见,我始终建议在使用完毕后拆除任何创建的 AWS 基础设施。如果使用 CloudFormation,可以通过以下 AWS CLI 命令轻松完成:

code
aws cloudformation delete-stack \
  --region us-east-2 \
  --stack-name cluster-duck-test-v2

总结

我创建了 "cluster-duck" 仓库来测试 DuckDB 新的 Quack 通信协议。Quack 允许不同服务器上的 DuckDB 数据库通过 HTTP 相互"通信",DuckDB 将其定位为 DuckDB 数据库之间客户端/服务器通信的使能者。

这项发展可能非常有用,我特别感兴趣的是观察 Quack 在并发读写远程数据库时的表现。

在我的测试中,以及我演示的示例中,答案似乎是它处理得相当不错。

DuckDB 团队表示 Quack 是一个实验性功能,目前仍在开发中。正因如此,协议、函数名称、设置和默认值可能会发生变更,因此请务必不要将 Quack 用于任何生产系统。

希望本文中阐述的概念能对您有所帮助,如果您需要在运行于不同服务器上的 DuckDB 数据库上并行执行查询或其他 SQL 语句,这些内容将非常有用。

我不禁好奇 DuckDB 对于 Quack 的未来规划。如果 Quack 成为 DuckDB 生态系统中得到全面支持的一部分,我完全可以预见 Quack 可能会为未来的分布式 DuckDB 引擎提供传输和会话处理功能,但这种情况可能还很遥远。即便 Quack 只实现其当前具备的功能,它本身也具有独立的实用价值。

到此为止,我就先介绍到这里。您可以通过以下链接访问我的 GitHub 仓库,获取所有代码、CloudFormation 模板等资源:

https://github.com/taupirho/cluster-duck

有关 DuckDB 和 Quack 的更多信息,您可以通过以下链接查看 DuckDB 官方文档:

https://duckdb.org/docs/current

PS:目前我正在寻找合同工作机会。如果您或您认识的人正在寻找一位经验丰富的数据工程师,无论远程工作还是基于英国爱丁堡的职位,要求具备 AWS、AI、Python、SQL、PySpark、DuckDB 等技能,请通过 LinkedIn 与我联系。

作者:Thomas Reid

查看 Thomas Reid 的更多文章

数据科学

Duckdb

编程

Python

关系型数据库

分享本文

  • 在 Facebook 上分享
  • 在 LinkedIn 上分享
  • 在 X 上分享

Towards Data Science 是一个社区出版物。提交您的见解以触达全球受众,并通过 TDS 作者支付计划获得报酬。

请将 href 更新为您的实际投稿链接

为 TDS 写作

✦ 结束 CTA ✦