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

如何用C# UDP监听器持续接收UDP推送的JSON定位数据

C# UDP持续接收定位JSON数据并入库解决方案

问题根源

你当前的代码只调用了一次QuuppaListener.Receive(),这个方法是阻塞式的,执行一次仅接收一条数据。要实现持续接收,需要将接收逻辑放入循环中,同时结合多线程/异步处理大量数据的解析和入库操作,避免阻塞监听流程。

改造后的完整代码

首先定义与JSON结构匹配的实体类,方便解析数据:

public class TagLocationData
{
    public string tagId { get; set; }
    public string tagName { get; set; }
    public long lastPacketTS { get; set; }
    public string color { get; set; }
    public string tagGroupName { get; set; }
    public string locationType { get; set; }
    public string locationMovementStatus { get; set; }
    public double locationRadius { get; set; }
    public double[] location { get; set; }
    public long locationTS { get; set; }
    public string locationCoordSysId { get; set; }
    public string locationCoordSysName { get; set; }
    public string[] locationZoneIds { get; set; }
    public string[] locationZoneNames { get; set; }
    public string button1State { get; set; }
    public long button1StateTS { get; set; }
    public long? button1LastPressTS { get; set; }
    public string batteryAlarm { get; set; }
    public long batteryAlarmTS { get; set; }
    public int rssi { get; set; }
    public int rssiLocatorCount { get; set; }
    public long lastSeenTS { get; set; }
}

主程序改造为持续监听、批量入库的版本:

using System;
using System.Net;
using System.Net.Sockets;
using System.Text;
using System.Text.Json;
using System.Threading;
using System.Data.SqlClient;
using System.Collections.Generic;
using System.IO;

namespace QuuppaUDPListener
{
    class Program
    {
        // 替换为你的数据库连接字符串
        private const string DbConnectionString = "Server=你的服务器地址;Database=你的数据库名;User Id=用户名;Password=密码;";
        private static UdpClient _udpListener;
        private static readonly List<TagLocationData> _pendingData = new List<TagLocationData>();
        // 批量插入阈值:达到该数量自动触发入库
        private const int BatchInsertThreshold = 1000;
        private static readonly object _locker = new object();

        static void Main(string[] args)
        {
            _udpListener = new UdpClient(6900);
            var remoteEndpoint = new IPEndPoint(IPAddress.Any, 6900);
            
            Console.WriteLine("UDP监听器已启动,等待定位数据...");

            // 启动后台线程处理定时批量入库,避免阻塞接收流程
            var dbProcessingThread = new Thread(ProcessTimedBatchInsert) { IsBackground = true };
            dbProcessingThread.Start();

            try
            {
                // 无限循环持续接收数据
                while (true)
                {
                    byte[] receivedBytes = _udpListener.Receive(ref remoteEndpoint);
                    string rawJson = Encoding.UTF8.GetString(receivedBytes, 0, receivedBytes.Length);
                    
                    // 解析JSON数据
                    try
                    {
                        var locationRecord = JsonSerializer.Deserialize<TagLocationData>(rawJson);
                        if (locationRecord != null)
                        {
                            lock (_locker)
                            {
                                _pendingData.Add(locationRecord);
                                // 达到批量阈值时触发异步入库
                                if (_pendingData.Count >= BatchInsertThreshold)
                                {
                                    var batchToInsert = new List<TagLocationData>(_pendingData);
                                    _pendingData.Clear();
                                    Task.Run(() => InsertBatchToDatabase(batchToInsert));
                                }
                            }
                            Console.WriteLine($"已接收数据:TagID={locationRecord.tagId}");
                        }
                    }
                    catch (JsonException jsonEx)
                    {
                        Console.WriteLine($"JSON解析失败:{jsonEx.Message},原始数据:{rawJson}");
                    }
                }
            }
            catch (SocketException sockEx)
            {
                Console.WriteLine($"Socket异常:{sockEx.Message}");
            }
            catch (Exception ex)
            {
                Console.WriteLine($"未知异常:{ex.Message}");
            }
            finally
            {
                _udpListener.Close();
                // 程序退出前插入剩余数据
                lock (_locker)
                {
                    if (_pendingData.Count > 0)
                    {
                        InsertBatchToDatabase(_pendingData);
                    }
                }
            }
        }

        /// <summary>
        /// 定时检查并插入剩余数据(每20秒执行一次,匹配数据源推送间隔)
        /// </summary>
        private static void ProcessTimedBatchInsert()
        {
            while (true)
            {
                Thread.Sleep(20000);
                lock (_locker)
                {
                    if (_pendingData.Count > 0)
                    {
                        var batchToInsert = new List<TagLocationData>(_pendingData);
                        _pendingData.Clear();
                        InsertBatchToDatabase(batchToInsert);
                    }
                }
            }
        }

        /// <summary>
        /// 批量插入数据到TagHistoryTbl表
        /// </summary>
        private static void InsertBatchToDatabase(List<TagLocationData> dataBatch)
        {
            if (dataBatch.Count == 0) return;

            try
            {
                using (var dbConn = new SqlConnection(DbConnectionString))
                {
                    dbConn.Open();
                    // 构建批量插入SQL(也可使用SqlBulkCopy提升效率和安全性)
                    var sqlBuilder = new StringBuilder();
                    sqlBuilder.Append(@"INSERT INTO TagHistoryTbl 
(tagId, tagName, lastPacketTS, color, tagGroupName, locationType, locationMovementStatus, locationRadius, 
locationX, locationY, locationZ, locationTS, locationCoordSysId, locationCoordSysName, locationZoneIds, 
locationZoneNames, button1State, button1StateTS, button1LastPressTS, batteryAlarm, batteryAlarmTS, rssi, 
rssiLocatorCount, lastSeenTS) VALUES ");

                    for (int i = 0; i < dataBatch.Count; i++)
                    {
                        var record = dataBatch[i];
                        // 将数组类型字段序列化为JSON字符串存储
                        string zoneIdsJson = JsonSerializer.Serialize(record.locationZoneIds);
                        string zoneNamesJson = JsonSerializer.Serialize(record.locationZoneNames);

                        sqlBuilder.Append($"('{record.tagId}', '{record.tagName}', {record.lastPacketTS}, '{record.color}', " +
                            $"'{record.tagGroupName}', '{record.locationType}', '{record.locationMovementStatus}', {record.locationRadius}, " +
                            $"{record.location[0]}, {record.location[1]}, {record.location[2]}, {record.locationTS}, " +
                            $"'{record.locationCoordSysId}', '{record.locationCoordSysName}', '{zoneIdsJson}', '{zoneNamesJson}', " +
                            $"'{record.button1State}', {record.button1StateTS}, {(record.button1LastPressTS.HasValue ? record.button1LastPressTS.Value.ToString() : "NULL")}, " +
                            $"'{record.batteryAlarm}', {record.batteryAlarmTS}, {record.rssi}, {record.rssiLocatorCount}, {record.lastSeenTS})");

                        if (i < dataBatch.Count - 1)
                        {
                            sqlBuilder.Append(",");
                        }
                    }

                    using (var cmd = new SqlCommand(sqlBuilder.ToString(), dbConn))
                    {
                        int insertedRows = cmd.ExecuteNonQuery();
                        Console.WriteLine($"批量入库完成:共插入{insertedRows}条记录");
                    }
                }
            }
            catch (SqlException sqlEx)
            {
                Console.WriteLine($"数据库操作失败:{sqlEx.Message}");
                // 将失败数据写入本地文件,便于后续重试
                File.AppendAllText("db_insert_errors.log", $"[{DateTime.Now:yyyy-MM-dd HH:mm:ss}] 错误信息:{sqlEx.Message}\r\n失败数据:{JsonSerializer.Serialize(dataBatch)}\r\n");
            }
            catch (Exception ex)
            {
                Console.WriteLine($"批量处理失败:{ex.Message}");
            }
        }
    }
}

关键优化点

  • 持续监听:通过while(true)循环调用Receive(),实现不间断接收UDP数据包
  • 批量入库:设置阈值(1000条)+定时(20秒)双机制,避免频繁操作数据库,提升性能
  • 线程安全:用lock保护共享数据列表,防止多线程冲突
  • 错误处理:分别捕获Socket、JSON解析、数据库异常,记录错误信息并保留失败数据
  • 异步处理:后台线程处理数据库操作,不会阻塞UDP接收流程

注意事项

  1. 确保TagHistoryTbl表结构与实体类字段匹配,数组类型字段可选择拆分存储或存为JSON字符串
  2. 替换代码中的数据库连接字符串为实际环境信息
  3. 若使用.NET Framework,需安装Newtonsoft.Json NuGet包替代System.Text.Json
  4. 数据量极大时,建议使用SqlBulkCopy替代拼接SQL的方式,更安全高效
  5. 可添加程序退出信号监听(如Console.ReadLine()捕获退出指令),确保剩余数据入库

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 05:07:05