Running SQL Concurrently Across Three Remote DuckDB Servers with Quack
TL;DR · AI 摘要
Quack协议实现跨三台远程DuckDB服务器的SQL并发执行,通过AWS EC2实验验证可行性。
核心要点
- Quack协议通过HTTP实现远程DuckDB数据库通信,支持跨服务器读写操作。
- 实验部署3台AWS EC2服务器,使用ARM64架构和DuckDB 1.5.5版本。
- 并发SQL执行需配置Quack扩展、systemd服务及自动关机计时器。
结构提纲
按章节快速跳转。
- §引言
介绍Quack协议的远程数据库通信能力及实验目标。
通过CloudFormation部署3台AWS EC2服务器并安装依赖组件。
详细说明Quack扩展、服务配置及认证令牌设置方法。
展示如何通过协调节点并行执行跨服务器SQL查询。
- ›结果分析
对比单节点与多节点执行性能差异及潜在优化方向。
思维导图
用一张图看清主题之间的关系。
查看大纲文本(无障碍 / 无 JS 友好)
- Quack远程SQL并发实验
- 核心组件
- DuckDB 1.5.5
- Quack协议
- AWS EC2
- 部署流程
- CloudFormation配置
- Python环境搭建
- systemd服务设置
- 实验验证
- 并发SQL执行
- 性能对比
- 自动关机机制
金句 / Highlights
值得收藏与分享的关键句。
Quack协议允许DuckDB在远程服务器间进行读写操作,但不支持分布式查询处理。
实验环境使用Amazon Linux 2023 ARM64和Python 3.12虚拟环境部署。
Quack默认监听9494端口,需配置认证令牌实现安全通信。
通过 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,也可以在本地)部署堆栈,命令如下:
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。
这些文件安装在每台服务器上:
/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)还会获得这些命令启动器:
/usr/local/bin/cluster-duck
/usr/local/bin/cluster-duck-sql在协调器上运行 cluster-duck-sql 时,最终会启动:
/opt/cluster-duck-venv/bin/python
/opt/cluster-duck/related_cluster_sql.pyPython 源代码被压缩并直接嵌入到 CloudFormation 模板中,作为 Base64 编码的存档。
在 EC2 引导过程中,用户数据脚本会执行以下操作:
- 解码嵌入的存档。
- 创建 /opt/cluster-duck 目录。
- 将 Python 文件提取到该目录中。
- 在 worker 1 上创建命令启动器。
- 在每个 worker 上启动 Quack 服务。
所有需要的内容都包含在 CloudFormation 文件中。
协调机制的工作原理
Quack 会将每个 SQL 语句发送到选定的 DuckDB 服务器,并返回结果。协调三个调用的部分是运行在 Worker 1 上的 Python 代码。
首先,每个 —query 或 —query-file 参数都会被验证为单个 SQL 语句,并转换为 QueryFragment。这些片段会按照提供的顺序进行标记:
fragments.append(
QueryFragment(worker_id, f"query-{index}", sql)
)然后协调器为每个片段创建一个线程,并创建一个参与人数相同的屏障:
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() 发送语句:
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 数据库配置如下。
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条记录,帮助您更好地了解表中内容。
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 标签:
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 在每台服务器上安装了这个程序:
/opt/cluster-duck/seed_related_data.py然后执行:
/opt/cluster-duck-venv/bin/python /opt/cluster-duck/seed_related_data.py \
--worker "$WORKER_INDEX" \
--rows 10000000行数来自 CloudFormation 的 RowCount 参数,该参数默认值为 1000 万。
3. 打开工作节点的 DuckDB 文件
该程序打开:
/var/lib/cluster-duck/worker.duckdb其中包含以下代码:
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 存储在单独的文件中并使用这些文件作为输入。例如:
[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 中,输入以下代码:
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
### 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
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不会回滚成功的语句。
为了使演示可重复,首先删除任何早期运行留下的行:
[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参数。
[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个新表,然后查询它们。
[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 认证令牌。其工作流程如下:
创建 CloudFormation 堆栈
↓
调用 TokenManager Lambda
↓
生成三个随机的 64 字符令牌
↓
将它们存储为 SSM 中的 SecureString 参数
↓
返回成功并停止如果我们运行这个完整设置 4 小时,预估成本如下:
组件 预估成本
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 命令轻松完成:
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 ✦