multiTableRepartitionDS
语法
multiTableRepartitionDS(queries, [column], [partitionType],
[partitionScheme], [local=true])
详情
参数
queries 是一个 SQL 查询元代码,或由多个 SQL 查询元代码组成的元组。指定元组时,每个查询生成的数据源数量必须相同。指定单个元代码时必须同时指定 column;若只需按表原有分区生成数据源,应使用 sqlDS。
column 字符串标量,表示用于重新划分数据源的列名。输入为单个 SQL 查询元代码时必须指定;输入为元代码元组时可选。multiTableRepartitionDS 会根据该列划分数据源。
partitionType 可选参数,表示分区类型,可取 RANGE 或 VALUE。3.00 系列还支持 HASH。
partitionScheme 可选参数,是一个向量,表示分区方案。3.00 系列也支持标量。
local 可选参数,是一个布尔值,表示是否将数据源获取到当前节点进行计算。默认值为 true。
返回值
返回一个元组,包含一组数据源。
若 queries 是元代码元组,则返回嵌套元组,外层元组中的每个元素均为子元组,包含各查询在相同分区索引的数据源,子元组长度等于 queries 的长度。
例子
以下例子使用同一数据库中的订单表和成交表。
createCatalog("factor_pipeline")
go
use catalog factor_pipeline
create database market_data partitioned by VALUE(2026.01.02 2026.01.03), engine='OLAP'
go
create table market_data.orders (
TradeDate DATE,
SecurityID SYMBOL,
TradeTime TIME,
OrderNO LONG,
OrderQty INT
)
partitioned by TradeDate
go
create table market_data.trades (
TradeDate DATE,
SecurityID SYMBOL,
TradeTime TIME,
TradePrice DOUBLE,
OfferApplSeqNum LONG
)
partitioned by TradeDate
go
orderData = table(
2026.01.02 2026.01.02 2026.01.02 2026.01.03 2026.01.03 2026.01.03 as TradeDate,
`AAPL`AAPL`MSFT`AAPL`AAPL`MSFT as SecurityID,
09:30:00.000 09:31:00.000 09:32:00.000 09:30:00.000 09:31:00.000 09:32:00.000 as TradeTime,
1001 1002 1001 2001 2002 2001 as OrderNO,
100 300 200 150 150 400 as OrderQty
)
tradeData = table(
2026.01.02 2026.01.02 2026.01.02 2026.01.03 2026.01.03 2026.01.03 as TradeDate,
`AAPL`AAPL`MSFT`AAPL`AAPL`MSFT as SecurityID,
09:30:00.300 09:31:00.700 09:32:00.200 09:30:00.100 09:31:00.400 09:32:00.100 as TradeTime,
0.0 0.0 0.0 0.0 0.0 12.0 as TradePrice,
1001 1002 1001 2001 2002 2001 as OfferApplSeqNum
)
factor_pipeline.market_data.orders.append!(orderData)
factor_pipeline.market_data.trades.append!(tradeData)
例 1:根据表原有的分区方案生成对齐的数据源,不指定 column、partitionType 和 partitionScheme。
ds = multiTableRepartitionDS(queries=[
<select * from market_data.orders>,
<select * from market_data.trades>
])
ds.size()
// output: 2
ds[0].size()
// output: 2
例 2:根据股票代码的值划分数据源。
ds = multiTableRepartitionDS(
queries=[
<select * from market_data.orders>,
<select * from market_data.trades>
],
column=`SecurityID,
partitionType=VALUE,
partitionScheme=symbol(`AAPL`MSFT)
)
ds.size()
// output: 2
ds[0].size()
// output: 2
例 3:根据日期范围划分数据源。
ds = multiTableRepartitionDS(
queries=[
<select * from market_data.orders>,
<select * from market_data.trades>
],
column=`TradeDate,
partitionType=RANGE,
partitionScheme=2026.01.02 2026.01.03 2026.01.04
)
ds.size()
// output: 2
ds[0].size()
// output: 2
例 4:将按原分区方案生成的对齐数据源传给 mr。map 函数同时接收同一交易日的委托和成交数据,计算各股票
500 毫秒内撤单委托量占总委托量的比例。成交表中 TradePrice 为 0 的记录表示撤单。
ds = multiTableRepartitionDS(queries=[
<select * from market_data.orders>,
<select * from market_data.trades>
])
def calcCancelRatio(orderTB, tradeTB){
startTime = 09:30:00.000
endTime = 14:57:00.000
cancelTB = select TradeDate, SecurityID, TradeTime as CancelTime, OfferApplSeqNum as OrderNO
from tradeTB
where TradeTime between startTime and endTime, TradePrice=0
joinedTB = lj(orderTB, cancelTB, `TradeDate`SecurityID`OrderNO)
return select TradeDate, SecurityID,
1.0 * sum(iif(isNull(CancelTime), 0,
iif((CancelTime-TradeTime>=0) and (CancelTime-TradeTime<500), OrderQty, 0))) /
sum(OrderQty) as cancelRatio
from joinedTB
where TradeTime between startTime and endTime
group by TradeDate, SecurityID
}
result = mr(ds=ds, mapFunc=calcCancelRatio, reduceFunc=unionAll)
result.sortBy!(`TradeDate`SecurityID)
result
输出如下:
| TradeDate | SecurityID | cancelRatio |
|---|---|---|
| 2026.01.02 | AAPL | 0.25 |
| 2026.01.02 | MSFT | 1 |
| 2026.01.03 | AAPL | 1 |
| 2026.01.03 | MSFT | 0 |
相关函数:repartitionDS、mr
