You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在Azure Stream Analytics中基于最近500交易单位计算成交量加权平均价

Azure Stream Analytics实现基于最近500交易单位的VWAP计算

你需要按产品分组,基于最近500交易单位计算成交量加权平均价(VWAP),而非依赖时间窗口或全量数据聚合。Azure Stream Analytics(AS)原生没有直接的“数量窗口”函数,但这个需求完全可行,核心思路是通过**用户定义聚合(UDA)**维护每个产品的交易状态,动态筛选出累计数量不超过500的最近交易,再实时计算VWAP。

实现方案:用户定义聚合(UDA)

VWAP的计算公式为:(Σ(price × quantity)) / Σ(quantity),我们需要让UDA自动维护每个产品的交易队列,当累计交易单位超过500时,移除最早的交易记录,始终保留最近的、总数量≤500的交易,再实时计算VWAP。

1. 创建JavaScript UDA

在Azure Stream Analytics作业中添加JavaScript类型的用户定义聚合,代码如下:

function main() {
    this.init = function () {
        this.transactionQueue = [];
        this.totalQuantity = 0;
        this.totalValue = 0;
    }

    this.accumulate = function (transaction) {
        // 新增交易到队列
        this.transactionQueue.push({
            price: transaction.price,
            qty: transaction.quantity
        });
        this.totalQuantity += transaction.quantity;
        this.totalValue += transaction.price * transaction.quantity;

        // 移除最早交易,确保总数量不超过500
        while (this.totalQuantity > 500) {
            const oldestTx = this.transactionQueue.shift();
            this.totalQuantity -= oldestTx.qty;
            this.totalValue -= oldestTx.price * oldestTx.qty;
        }
    }

    this.computeResult = function () {
        return this.totalQuantity === 0 ? 0 : this.totalValue / this.totalQuantity;
    }

    this.merge = function (otherState) {
        // 合并多分区的状态(处理作业并行场景)
        otherState.transactionQueue.forEach(tx => {
            this.transactionQueue.push(tx);
            this.totalQuantity += tx.qty;
            this.totalValue += tx.price * tx.qty;
        });
        // 合并后再次校验并裁剪队列
        while (this.totalQuantity > 500) {
            const oldestTx = this.transactionQueue.shift();
            this.totalQuantity -= oldestTx.qty;
            this.totalValue -= oldestTx.price * oldestTx.qty;
        }
    }
}

2. 在AS查询中调用UDA

使用滚动窗口触发聚合计算(窗口大小可根据业务需求调整,比如每秒触发一次),UDA内部会自动维护最近500交易单位的状态:

SELECT
    product,
    VWAPLast500Units(STRUCT(price, quantity)) AS vwap_last_500_units
INTO
    [你的输出目标(如Blob存储、SQL数据库等)]
FROM
    [你的Event Hub输入源]
GROUP BY
    product,
    TumblingWindow(second, 1)

关键说明

  • 这里的滚动窗口仅用于触发计算的频率,不限制数据的时间范围,UDA内部的状态完全基于交易数量维护,符合你“无时间约束、仅看最近500交易单位”的需求。
  • UDA支持多分区合并,适配AS作业的并行处理场景,确保每个产品的状态准确。

内容的提问来源于stack exchange,提问作者josh66

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 20:41:59