diff --git a/XCode/DataAccessLayer/Database/InfluxDB.cs b/XCode/DataAccessLayer/Database/InfluxDB.cs index 0c1dddfb1..d7d93e263 100644 --- a/XCode/DataAccessLayer/Database/InfluxDB.cs +++ b/XCode/DataAccessLayer/Database/InfluxDB.cs @@ -1,5 +1,6 @@ using System.Data; using System.Data.Common; +using System.Globalization; using System.Text; using NewLife.Collections; using NewLife.Data; @@ -51,6 +52,8 @@ public override Boolean Support(String providerName) #endregion #region 数据库特性 + public override BatchCapability BatchCapability => BatchCapability.Insert | BatchCapability.Upsert; + protected override String ReservedWordsStr => "AND,OR,NOT,FROM,WHERE,SELECT,DELETE,DROP,SHOW,MEASUREMENT,TAG,FIELD,TIME"; /// 格式化关键字 @@ -69,18 +72,30 @@ public override String FormatKeyWord(String keyWord) /// public override String FormatValue(IDataColumn field, Object? value) { + if (value == null) + return field.Nullable ? "null" : ""; + var code = System.Type.GetTypeCode(field.DataType); if (code == TypeCode.String) { - if (value == null) - return field.Nullable ? "null" : "\"\""; - return "\"" + value.ToString()?.Replace("\"", "\\\"") + "\""; } else if (code == TypeCode.Boolean) { return value.ToBoolean() ? "true" : "false"; } + else if (code is TypeCode.SByte or TypeCode.Byte or TypeCode.Int16 or TypeCode.UInt16 or TypeCode.Int32 or TypeCode.UInt32 or TypeCode.Int64) + { + return $"{value}i"; + } + else if (code == TypeCode.UInt64) + { + return $"{value}u"; + } + else if (code is TypeCode.Single or TypeCode.Double or TypeCode.Decimal) + { + return Convert.ToString(value, CultureInfo.InvariantCulture) ?? "0"; + } return base.FormatValue(field, value); } @@ -148,76 +163,81 @@ public override Task InsertAndGetIdentityAsync(String sql, CommandType ty /// 实体列表 /// public override Int32 Insert(IDataTable table, IDataColumn[] columns, IEnumerable list) + { + var lineProtocol = BuildLineProtocol(Database, table, columns, list); + return Execute(lineProtocol); + } + + /// 批量插入或更新 + /// 数据表 + /// 要插入的字段 + /// 主键已存在时,要更新的字段 + /// 主键已存在时,要累加更新的字段 + /// 实体列表 + /// + public override Int32 Upsert(IDataTable table, IDataColumn[] columns, ICollection? updateColumns, ICollection? addColumns, IEnumerable list) + { + // InfluxDB 自动处理相同 measurement + tags + timestamp 的写入,新值会覆盖旧值 + return Insert(table, columns, list); + } + + private static String BuildLineProtocol(IDatabase database, IDataTable table, IDataColumn[] columns, IEnumerable list) { var sb = Pool.StringBuilder.Get(); - var db = (Database as DbBase)!; + var db = (database as DbBase)!; - // InfluxDB 使用 Line Protocol 格式写入 - // 格式: measurement,tag1=value1,tag2=value2 field1=value1,field2=value2 timestamp foreach (var entity in list) { - // measurement 名称(表名) + var timeCol = columns.FirstOrDefault(c => + { + var name = c.Name ?? c.ColumnName; + return !name.IsNullOrEmpty() && name.EqualIgnoreCase("Time", "CreateTime", "UpdateTime") && entity[name] != null; + }); + sb.Append(db.FormatName(table)); - // tags(索引字段,通常是维度) - var tags = columns.Where(c => c.PrimaryKey || c.Master).ToArray(); + var tags = columns.Where(c => (c.PrimaryKey || c.Master) && c != timeCol).ToArray(); if (tags.Length > 0) { sb.Append(','); sb.Append(tags.Join(",", c => { - var value = entity[c.Name]; - return $"{db.FormatName(c)}={value}"; + var name = c.Name ?? c.ColumnName; + return $"{db.FormatName(c)}={entity[name]}"; })); } - // fields(数据字段) - var fields = columns.Where(c => !c.PrimaryKey && !c.Master).ToArray(); + var fields = columns.Where(c => !c.PrimaryKey && !c.Master && c != timeCol).ToArray(); if (fields.Length > 0) { sb.Append(' '); sb.Append(fields.Join(",", c => { - var value = entity[c.Name]; - var strValue = value?.ToString() ?? ""; - // 字符串字段需要加引号 - if (c.DataType == typeof(String)) - strValue = $"\"{strValue}\""; - return $"{db.FormatName(c)}={strValue}"; + var name = c.Name ?? c.ColumnName; + return $"{db.FormatName(c)}={db.FormatValue(c, entity[name])}"; })); } - // timestamp(纳秒级时间戳) - var timeCol = columns.FirstOrDefault(c => c.Name.EqualIgnoreCase("Time", "CreateTime", "UpdateTime")); if (timeCol != null) { - var time = entity[timeCol.Name]; + var name = timeCol.Name ?? timeCol.ColumnName; + var time = entity[name]; if (time is DateTime dt) - { - var epoch = new DateTime(1970, 1, 1, 0, 0, 0, DateTimeKind.Utc); - var nanos = (dt.ToUniversalTime() - epoch).Ticks * 100; - sb.Append($" {nanos}"); - } + sb.Append($" {ToNanoseconds(dt)}"); + else if (time is DateTimeOffset dto) + sb.Append($" {ToNanoseconds(dto.UtcDateTime)}"); } - sb.AppendLine(); + sb.Append('\n'); } - var lineProtocol = sb.Return(true); - return Execute(lineProtocol); + return sb.Return(true); } - /// 批量插入或更新 - /// 数据表 - /// 要插入的字段 - /// 主键已存在时,要更新的字段 - /// 主键已存在时,要累加更新的字段 - /// 实体列表 - /// - public override Int32 Upsert(IDataTable table, IDataColumn[] columns, ICollection? updateColumns, ICollection? addColumns, IEnumerable list) + private static Int64 ToNanoseconds(DateTime dt) { - // InfluxDB 自动处理相同 measurement + tags + timestamp 的写入,新值会覆盖旧值 - return Insert(table, columns, list); + var epoch = new DateTime(1970, 1, 1, 0, 0, 0, DateTimeKind.Utc); + return (dt.ToUniversalTime() - epoch).Ticks * 100; } #endregion diff --git a/XUnitTest.XCode/DataAccessLayer/BatchCapabilityTests.cs b/XUnitTest.XCode/DataAccessLayer/BatchCapabilityTests.cs index e3809f09e..c78f20f0a 100644 --- a/XUnitTest.XCode/DataAccessLayer/BatchCapabilityTests.cs +++ b/XUnitTest.XCode/DataAccessLayer/BatchCapabilityTests.cs @@ -103,6 +103,20 @@ public void NovaDb_BatchCapability() Assert.False(cap.HasFlag(BatchCapability.Update)); } + [Fact] + [System.ComponentModel.Description("InfluxDB具备Insert/Upsert能力,不含Update/InsertIgnore/Replace")] + public void InfluxDB_BatchCapability() + { + var db = DbFactory.Create(DatabaseType.InfluxDB); + var cap = db.BatchCapability; + + Assert.True(cap.HasFlag(BatchCapability.Insert)); + Assert.True(cap.HasFlag(BatchCapability.Upsert)); + Assert.False(cap.HasFlag(BatchCapability.Update)); + Assert.False(cap.HasFlag(BatchCapability.InsertIgnore)); + Assert.False(cap.HasFlag(BatchCapability.Replace)); + } + [Fact] [System.ComponentModel.Description("BatchCapability枚举值满足Flags语义,组合标志可通过HasFlag判断")] public void BatchCapability_FlagsSemantics() diff --git a/XUnitTest.XCode/DataAccessLayer/InfluxDBLineProtocolTests.cs b/XUnitTest.XCode/DataAccessLayer/InfluxDBLineProtocolTests.cs new file mode 100644 index 000000000..d41d584a3 --- /dev/null +++ b/XUnitTest.XCode/DataAccessLayer/InfluxDBLineProtocolTests.cs @@ -0,0 +1,88 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Reflection; +using NewLife.Data; +using XCode.DataAccessLayer; +using Xunit; + +namespace XUnitTest.XCode.DataAccessLayer; + +/// InfluxDB Line Protocol 纯单元测试,无需数据库连接 +public class InfluxDBLineProtocolTests +{ + [Fact] + public void BuildLineProtocol_ShouldFormatFieldTypes_AndSkipTimeField() + { + var db = DbFactory.Create(DatabaseType.InfluxDB); + var method = GetBuildLineProtocolMethod(); + + var table = DAL.CreateTable(); + table.TableName = "temperature"; + + var id = table.CreateColumn(); + id.ColumnName = "DeviceId"; + id.PrimaryKey = true; + id.DataType = typeof(Int32); + table.Columns.Add(id); + + var count = table.CreateColumn(); + count.ColumnName = "Count"; + count.DataType = typeof(Int32); + table.Columns.Add(count); + + var enabled = table.CreateColumn(); + enabled.ColumnName = "Enabled"; + enabled.DataType = typeof(Boolean); + table.Columns.Add(enabled); + + var name = table.CreateColumn(); + name.ColumnName = "Name"; + name.DataType = typeof(String); + table.Columns.Add(name); + + var time = table.CreateColumn(); + time.ColumnName = "Time"; + time.DataType = typeof(DateTime); + table.Columns.Add(time); + + var dt = new DateTime(2026, 7, 15, 0, 0, 0, DateTimeKind.Utc); + var model = new PlainModel + { + ["DeviceId"] = 1001, + ["Count"] = 7, + ["Enabled"] = true, + ["Name"] = "sensor \"A\"", + ["Time"] = dt + }; + + var lineProtocol = (String)method.Invoke(null, [db, table, table.Columns.ToArray(), new List { model }])!; + + Assert.Contains("Count=7i", lineProtocol); + Assert.Contains("Enabled=true", lineProtocol); + Assert.Contains("Name=\"sensor \\\"A\\\"\"", lineProtocol); + Assert.DoesNotContain("Time=", lineProtocol, StringComparison.OrdinalIgnoreCase); + Assert.DoesNotContain("\r\n", lineProtocol); + Assert.EndsWith("\n", lineProtocol); + + var nanos = (dt - new DateTime(1970, 1, 1, 0, 0, 0, DateTimeKind.Utc)).Ticks * 100; + Assert.Contains($" {nanos}\n", lineProtocol); + } + + private static MethodInfo GetBuildLineProtocolMethod() + { + var sessionType = typeof(DbFactory).Assembly.GetType("XCode.DataAccessLayer.InfluxDBSession", true)!; + return sessionType.GetMethod("BuildLineProtocol", BindingFlags.NonPublic | BindingFlags.Static)!; + } +} + +file class PlainModel : IModel +{ + private readonly Dictionary _data = new(StringComparer.OrdinalIgnoreCase); + + public Object? this[String name] + { + get => _data.GetValueOrDefault(name); + set => _data[name] = value; + } +}