[ PROMPT_NODE_24244 ]
R2 SQL 设计模式
[ SKILL_DOCUMENTATION ]
# R2 SQL 模式
R2 SQL 的常用模式、使用场景及集成示例。
## Wrangler CLI 查询
bash
# 基础查询
npx wrangler r2 sql query "my-bucket" "SELECT * FROM default.logs LIMIT 10"
# 多行查询
npx wrangler r2 sql query "my-bucket" "
SELECT status, COUNT(*), AVG(response_time)
FROM logs.http_requests
WHERE timestamp >= '2025-01-01T00:00:00Z'
GROUP BY status
ORDER BY COUNT(*) DESC
LIMIT 100
"
# 使用环境变量
export R2_SQL_WAREHOUSE="my-bucket"
npx wrangler r2 sql query "$R2_SQL_WAREHOUSE" "SELECT * FROM default.logs"
## HTTP API 查询
用于从外部系统进行程序化访问(非 Workers - 参见 gotchas.md)。
bash
curl -X POST https://api.cloudflare.com/client/v4/accounts/{account_id}/r2/sql/query
-H "Authorization: Bearer "
-H "Content-Type: application/json"
-d '{
"warehouse": "my-bucket",
"query": "SELECT * FROM default.my_table WHERE status = 200 LIMIT 100"
}'
响应:
{
"success": true,
"result": [{"user_id": "user_123", "timestamp": "2025-01-15T10:30:00Z", "status": 200}],
"errors": []
}
## Pipelines 集成
通过 Pipelines 将数据流式传输到 Iceberg 表,然后使用 R2 SQL 进行查询。
bash
# 设置 pipeline(选择 Data Catalog Table 作为目标)
npx wrangler pipelines setup
# 关键设置:
# - Destination: Data Catalog Table
# - Compression: zstd (推荐)
# - Roll file time: 300+ 秒 (生产环境), 10 秒 (开发环境)
# 发送数据到 pipeline
curl -X POST https://{stream-id}.ingest.cloudflare.com
-H "Content-Type: application/json"
-d '[{"user_id": "user_123", "event_type": "purchase", "timestamp": "2025-01-15T10:30:00Z", "amount": 29.99}]'
# 查询已摄入的数据 (等待滚动间隔)
npx wrangler r2 sql query "my-bucket" "
SELECT event_type, COUNT(*), SUM(amount)
FROM default.events
WHERE timestamp >= '2025-01-15T00:00:00Z'
GROUP BY event_type
"
详细设置请参阅 [pipelines/patterns.md](../pipelines/patterns.md)。
## PyIceberg 集成
使用 PyIceberg 创建并填充 Iceberg 表,然后使用 R2 SQL 进行查询。
python
from pyiceberg.catalog.rest import RestCatalog
import pyarrow as pa
import pandas as pd
# 设置 catalog
catalog = RestCatalog(
name="my_catalog",
warehouse="my-bucket",
uri="https://.r2.cloudflarestorage.com/iceberg/my-bucket",
token="",
)
catalog.create_namespace_if_not_exists("analytics")
# 创建表
schema = pa.schema([
pa.field