multiTableRepartitionDS

语法

multiTableRepartitionDS(queries, [column], [partitionType], [partitionScheme], [local=true])

详情

为多个表生成数据源,并可使用指定的分区类型和分区方案重新划分数据源。

queries 为 SQL 查询元代码组成的元组时,函数分别为每个查询生成数据源。若未指定 columnpartitionTypepartitionScheme,则对每个查询应用 sqlDS,沿用各表原有的分区方案。所有查询生成的数据源数量必须相同;函数按分区索引对齐,并将相同索引的数据源组成一个子元组。

若指定重分区参数,则对所有查询应用相同的分区方案,并按分区索引组织结果。将结果传给 mr 时,同一子元组中的数据源必须来自同一数据库的分布式表。

参数

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:根据表原有的分区方案生成对齐的数据源,不指定 columnpartitionTypepartitionScheme

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

相关函数:repartitionDSmr