修正命名错误
xiyunfei authored at 2023-04-09 21:10:58
1.87 KiB
NewLife.JT808
using JT808.Data;
using JT808.Worker.Workers;
using NewLife;
using NewLife.Caching;
using NewLife.Caching.Services;
using NewLife.Log;
using NewLife.Model;
using Stardust;
using Stardust.Extensions;
using XCode;
using XCode.DataAccessLayer;

XTrace.UseConsole();

var services = ObjectContainer.Current;

// 星尘注册
var star = services.AddStardust();

// Redis 分布式缓存
services.AddSingleton<ICacheProvider, RedisCacheProvider>();

// 数据库连接
var set = NewLife.Setting.Current;
if (set.IsNew)
{
    set.DataPath = "../Data";
    set.Save();
}

if (!DAL.ConnStrs.ContainsKey("GPS"))
    DAL.AddConnStr("GPS", "Server=127.0.0.1;Port=3306;Database=JT808;User=root;Password=root;", null, "MySql");
if (!DAL.ConnStrs.ContainsKey("Safety")) DAL.AddConnStr("Safety", "MapTo=GPS", null, "MySql");
if (!DAL.ConnStrs.ContainsKey("Platform")) DAL.AddConnStr("Platform", "MapTo=GPS", null, "MySql");

EntityFactory.InitConnection("GPS");
EntityFactory.InitConnection("Safety");
EntityFactory.InitConnection("Platform");

services.AddSingleton<QueueService>();

var host = services.BuildHost();

// 注册后台 Worker(每个 Worker 对应一个 Redis Stream 消费组)
host.Add<PositionWorker>();
host.Add<SensorWorker>();
host.Add<EventWorker>();
host.Add<ADASAlarmWorker>();
host.Add<DSMAlarmWorker>();
host.Add<BSDAlarmWorker>();
host.Add<TirePressureWorker>();

XTrace.WriteLine("JT808.Worker 已启动,共 7 个 Worker");

await host.RunAsync();

/// <summary>Redis Stream 消费者服务</summary>
public class QueueService
{
    private readonly ICacheProvider _cacheProvider;

    public QueueService(ICacheProvider cacheProvider)
    {
        _cacheProvider = cacheProvider;
    }

    public IProducerConsumer<T>? Create<T>(String topic, String? group = null)
    {
        var queue = _cacheProvider.GetQueue<T>(topic, group ?? "Group");
        return queue as IProducerConsumer<T>;
    }
}