Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
100 changes: 60 additions & 40 deletions XCode/DataAccessLayer/Database/InfluxDB.cs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
using System.Data;
using System.Data.Common;
using System.Globalization;
using System.Text;
using NewLife.Collections;
using NewLife.Data;
Expand Down Expand Up @@ -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";

/// <summary>格式化关键字</summary>
Expand All @@ -69,18 +72,30 @@ public override String FormatKeyWord(String keyWord)
/// <returns></returns>
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";
}
Comment on lines 73 to +98

return base.FormatValue(field, value);
}
Expand Down Expand Up @@ -148,76 +163,81 @@ public override Task<Int64> InsertAndGetIdentityAsync(String sql, CommandType ty
/// <param name="list">实体列表</param>
/// <returns></returns>
public override Int32 Insert(IDataTable table, IDataColumn[] columns, IEnumerable<IModel> list)
{
var lineProtocol = BuildLineProtocol(Database, table, columns, list);
return Execute(lineProtocol);
}

/// <summary>批量插入或更新</summary>
/// <param name="table">数据表</param>
/// <param name="columns">要插入的字段</param>
/// <param name="updateColumns">主键已存在时,要更新的字段</param>
/// <param name="addColumns">主键已存在时,要累加更新的字段</param>
/// <param name="list">实体列表</param>
/// <returns></returns>
public override Int32 Upsert(IDataTable table, IDataColumn[] columns, ICollection<String>? updateColumns, ICollection<String>? addColumns, IEnumerable<IModel> list)
{
// InfluxDB 自动处理相同 measurement + tags + timestamp 的写入,新值会覆盖旧值
return Insert(table, columns, list);
}

private static String BuildLineProtocol(IDatabase database, IDataTable table, IDataColumn[] columns, IEnumerable<IModel> 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;
});
Comment on lines +191 to +195

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);
}

/// <summary>批量插入或更新</summary>
/// <param name="table">数据表</param>
/// <param name="columns">要插入的字段</param>
/// <param name="updateColumns">主键已存在时,要更新的字段</param>
/// <param name="addColumns">主键已存在时,要累加更新的字段</param>
/// <param name="list">实体列表</param>
/// <returns></returns>
public override Int32 Upsert(IDataTable table, IDataColumn[] columns, ICollection<String>? updateColumns, ICollection<String>? addColumns, IEnumerable<IModel> 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

Expand Down
14 changes: 14 additions & 0 deletions XUnitTest.XCode/DataAccessLayer/BatchCapabilityTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
88 changes: 88 additions & 0 deletions XUnitTest.XCode/DataAccessLayer/InfluxDBLineProtocolTests.cs
Original file line number Diff line number Diff line change
@@ -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;

/// <summary>InfluxDB Line Protocol 纯单元测试,无需数据库连接</summary>
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<IModel> { 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<String, Object?> _data = new(StringComparer.OrdinalIgnoreCase);

public Object? this[String name]
{
get => _data.GetValueOrDefault(name);
set => _data[name] = value;
}
}