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

如何通过Azure SQL触发器触发Azure Function并推送更新至SignalR流

问题描述

我正在基于SignalR开发实时数据网格,已搭建SignalR Hub,现有流程如下:

  • 客户端向服务器发送GetStocks消息
  • 服务器返回初始数据列表
  • 客户端订阅GetStockTickerStream流以接收更新

当前使用随机生成的股票数据,现需替换为真实数据,并实现当Azure SQL数据库中记录发生插入/更新/删除操作时,将更新推送到SignalR的GetStockTickerStream流。

请问是否可以为Azure SQL数据库设置触发器,使其在记录增删改时触发Azure Function?请提供详细实现说明及示例代码。


实现方案与步骤

完全可以实现,推荐采用Azure SQL变更数据捕获(CDC) + Azure SQL触发器类型的Function + SignalR服务的架构,能可靠捕获数据库变更并推送到SignalR客户端。以下是详细实现步骤和代码修改:

1. 准备Azure SQL数据库

1.1 创建Stock表

在Azure SQL中创建对应业务实体的表:

CREATE TABLE Stocks (
    Id BIGINT PRIMARY KEY IDENTITY(1,1),
    Symbol NVARCHAR(50) NOT NULL,
    Price FLOAT NOT NULL,
    UpdatedAt DATETIME DEFAULT GETDATE()
);

1.2 启用变更数据捕获(CDC)

为Stocks表启用CDC,确保能捕获所有增删改操作:

-- 启用数据库级CDC
EXEC sys.sp_cdc_enable_db;

-- 启用Stocks表的CDC,捕获全量变更
EXEC sys.sp_cdc_enable_table
    @source_schema = N'dbo',
    @source_name = N'Stocks',
    @role_name = NULL,
    @supports_net_changes = 1;

2. 创建Azure SQL触发器Function

创建Azure Function并选择SQL触发器类型,配置连接字符串指向你的Azure SQL数据库,触发器会自动捕获CDC中的变更记录。

2.1 Function核心代码(C#)

using Microsoft.Azure.WebJobs;
using Microsoft.Azure.WebJobs.Extensions.SignalRService;
using Microsoft.Extensions.Logging;

namespace StockUpdatesFunction
{
    public static class StockChangeTrigger
    {
        [FunctionName("StockChangeTrigger")]
        public static async Task Run(
            [SqlTrigger("[dbo].[Stocks]", ConnectionStringSetting = "AzureSqlConnection")] IReadOnlyList<SqlChange<Stock>> changes,
            [SignalR(HubName = "StockTickerHub")] IAsyncCollector<SignalRMessage> signalRMessages,
            ILogger log)
        {
            foreach (var change in changes)
            {
                log.LogInformation($"捕获到股票变更: {change.Operation} - ID:{change.Item.Id}");
                
                // 根据操作类型推送对应消息
                switch (change.Operation)
                {
                    case SqlOperation.Insert:
                    case SqlOperation.Update:
                        // 推送更新后的股票数据到SignalR
                        await signalRMessages.AddAsync(new SignalRMessage
                        {
                            Target = "GetStockTickerStream",
                            Arguments = new[] { change.Item }
                        });
                        break;
                    case SqlOperation.Delete:
                        // 推送删除通知,客户端自行处理移除逻辑
                        await signalRMessages.AddAsync(new SignalRMessage
                        {
                            Target = "StockDeleted",
                            Arguments = new[] { change.Item.Id }
                        });
                        break;
                }
            }
        }
    }

    // 与SQL表映射的实体类
    public class Stock
    {
        public long Id { get; set; }
        public string Symbol { get; set; }
        public double Price { get; set; }
    }
}

2.2 Function配置文件

在local.settings.json中添加连接字符串配置:

{
  "IsEncrypted": false,
  "Values": {
    "AzureWebJobsStorage": "UseDevelopmentStorage=true",
    "FUNCTIONS_WORKER_RUNTIME": "dotnet",
    "AzureSqlConnection": "Server=tcp:{你的SQL服务器}.database.windows.net,1433;Initial Catalog={你的数据库名};Persist Security Info=False;User ID={账号};Password={密码};MultipleActiveResultSets=False;Encrypt=True;TrustServerCertificate=False;Connection Timeout=30;",
    "AzureSignalRConnectionString": "Endpoint=https://{你的SignalR服务}.service.signalr.net;AccessKey={你的访问密钥};Version=1.0;"
  }
}

3. 修改SignalR服务端代码

3.1 更新StockTickerService,替换随机数据为SQL读取

using Microsoft.Data.SqlClient;
using System.Data;

public sealed class StockTickerService : IStockTickerService
{
    private readonly string _connectionString;

    public StockTickerService(IConfiguration configuration)
    {
        _connectionString = configuration.GetConnectionString("AzureSqlConnection");
    }

    public async Task<IEnumerable<Stock>> GetStocksAsync()
    {
        var stocks = new List<Stock>();
        using var conn = new SqlConnection(_connectionString);
        await conn.OpenAsync();
        var cmd = new SqlCommand("SELECT Id, Symbol, Price FROM Stocks", conn);
        using var reader = await cmd.ExecuteReaderAsync();
        while (await reader.ReadAsync())
        {
            stocks.Add(new Stock
            {
                Id = reader.GetInt64(0),
                Symbol = reader.GetString(1),
                Price = reader.GetDouble(2)
            });
        }
        return stocks;
    }

    // 移除原定时生成随机数据的逻辑,改为由Azure Function触发推送
    public Task SubscribeAsync(ChannelWriter<Stock> writer, CancellationToken cancellationToken)
    {
        // 流模式下保持Channel活跃,或改为直接通过HubContext推送
        return Task.CompletedTask;
    }
}

3.2 调整SignalR Hub,支持外部推送

public sealed class StockTickerHub : Hub
{
    private readonly IStockTickerService _stockTickerService;

    public StockTickerHub(IStockTickerService stockTickerService)
    {
        _stockTickerService = stockTickerService;
    }

    public async Task GetStocks()
    {
        var stocks = await _stockTickerService.GetStocksAsync();
        await Clients.Caller.SendAsync("ReceiveStocks", stocks);
    }

    // 供Azure Function调用的推送方法
    public async Task PushStockUpdate(Stock stock)
    {
        await Clients.All.SendAsync("GetStockTickerStream", stock);
    }
}

3.3 更新服务注册

在Program.cs中添加必要的服务配置:

builder.Services.AddDataGrid();
builder.Services.AddSignalR();
builder.Configuration.AddJsonFile("appsettings.json", optional: false);

4. 客户端代码调整

4.1 新增删除事件监听

import { useEffect } from "react";
import { useDispatch, useSelector } from "react-redux";
import { Table } from "antd";
import { HubConnectionState } from "redux-signalr";
import hubConnection from "../store/middlewares/signalr/signalrSlice";
import { Stock, addStock, removeStock } from "../store/reducers/stockSlice";
import { RootState } from "../store";
import "./DataGrid.css";

const DataGrid = () => {
  const dispatch = useDispatch();
  const stocks = useSelector((state: RootState) => state.stock.stocks);

  useEffect(() => {
    if (hubConnection.state !== HubConnectionState.Connected) {
      hubConnection
        .start()
        .then(() => {
          console.log("SignalR连接成功");

          hubConnection.send("GetStocks");

          // 监听股票更新流
          hubConnection.stream("GetStockTickerStream").subscribe({
            next: async (item: Stock) => {
              dispatch(addStock(item));
            },
            error: (err) => {
              console.error("流订阅错误:", err);
            },
          });

          // 监听股票删除事件
          hubConnection.on("StockDeleted", (stockId: number) => {
            dispatch(removeStock(stockId));
          });
        })
        .catch((err) => console.error("SignalR连接失败:", err.toString()));
    }
  }, [dispatch]);

  return (
    <Table dataSource={stocks} rowKey={(record) => record.id}>
      <Table.Column title="ID" dataIndex="id" key="id" />
      <Table.Column title="股票代码" dataIndex="symbol" key="symbol" />
      <Table.Column title="价格" dataIndex="price" key="price" />
    </Table>
  );
};

export default DataGrid;

4.2 更新stockSlice添加删除逻辑

import { createSlice, PayloadAction } from "@reduxjs/toolkit";

export type Stock = Readonly<{
  id: number;
  symbol: string;
  price: number;
}>;

export type StockState = Readonly<{
  stocks: Stock[];
}>;

const initialState: StockState = {
  stocks: [],
};

const stockSlice = createSlice({
  name: "stock",
  initialState: initialState,
  reducers: {
    getStocks: (state, action: PayloadAction<Stock[]>) => {
      state.stocks = action.payload; // 替换为全量数据,避免重复
    },
    addStock: (state, action: PayloadAction<Stock>) => {
      const stockIndex = state.stocks.findIndex(
        (stock) => stock.id === action.payload.id
      );
    
      if (stockIndex !== -1) {
        state.stocks[stockIndex].price = action.payload.price;
      } else {
        state.stocks.push(action.payload);
      }
    },
    removeStock: (state, action: PayloadAction<number>) => {
      state.stocks = state.stocks.filter(stock => stock.id !== action.payload);
    }
  },
});

export const { getStocks, addStock, removeStock } = stockSlice.actions;

export default stockSlice.reducer;

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 05:48:12