如何通过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
相关产品推荐
相关产品推荐

