加入若干pgsql的接口

This commit is contained in:
2025-12-13 16:33:41 +08:00
parent 97337f31ce
commit f9db3c7226
10 changed files with 1526 additions and 131 deletions

View File

@@ -1,171 +1,380 @@
using Npgsql;
using System.Diagnostics;
using Newtonsoft.Json.Linq;
using Npgsql;
using System;
using System.Collections.Generic;
using System.Data;
using System.Data.SQLite;
using System.Linq;
namespace Ramitta
{
public partial class PostgreSql
public class PostgreSQL : IDisposable
{
private readonly string _connectionString;
private NpgsqlConnection db;
private bool disposed = false;
public PostgreSql(string connectionString)
// 构造函数,初始化数据库连接
public PostgreSQL(string connectionString, bool readOnly = false)
{
_connectionString = connectionString;
var connString = connectionString;
if (readOnly)
{
if (!connString.Contains("CommandTimeout"))
{
connString += ";CommandTimeout=0";
}
}
db = new NpgsqlConnection(connString);
db.Open();
}
// 执行查询操作,返回查询结果
public async Task<List<Dictionary<string, object>>> ExecuteQueryAsync(string query, Dictionary<string, object> parameters = null)
// 创建表:根据表名和字段定义创建表
public void CreateTable(string tableName, Dictionary<string, string> columns)
{
// 构建列定义的字符串
var columnsDefinition = string.Join(", ", columns.Select(c => $"\"{c.Key}\" {MapDataType(c.Value)}"));
string createTableQuery = $"CREATE TABLE IF NOT EXISTS \"{tableName}\" ({columnsDefinition});";
using (var cmd = new NpgsqlCommand(createTableQuery, db))
{
cmd.ExecuteNonQuery();
}
}
// 数据类型映射SQLite到PostgreSQL
private string MapDataType(string sqliteType)
{
return sqliteType.ToUpper() switch
{
"INTEGER" => "INTEGER",
"TEXT" => "TEXT",
"REAL" => "REAL",
"BLOB" => "BYTEA",
"NUMERIC" => "NUMERIC",
"BOOLEAN" => "BOOLEAN",
"DATE" => "DATE",
"DATETIME" => "TIMESTAMP",
"TIMESTAMP" => "TIMESTAMP",
_ => "TEXT" // 默认映射为TEXT
};
}
// 获取数据库中所有表名
public List<string> GetAllTableNames()
{
List<string> tableNames = new List<string>();
string query = "SELECT table_name FROM information_schema.tables WHERE table_schema = 'public';";
using (var cmd = new NpgsqlCommand(query, db))
using (var reader = cmd.ExecuteReader())
{
while (reader.Read())
{
string tableName = reader.GetString(0);
tableNames.Add(tableName);
}
}
return tableNames;
}
// 向已存在的表中添加新列
public void AddColumn(string tableName, string columnName, string columnType)
{
// 检查表是否存在
if (!TableExists(tableName))
{
throw new ArgumentException($"表 '{tableName}' 不存在");
}
// 检查列是否已存在
if (ColumnExists(tableName, columnName))
{
Console.WriteLine($"列 '{columnName}' 在表 '{tableName}' 中已存在");
return;
}
// 构建添加列的SQL语句
string addColumnQuery = $"ALTER TABLE \"{tableName}\" ADD COLUMN \"{columnName}\" {MapDataType(columnType)};";
using (var cmd = new NpgsqlCommand(addColumnQuery, db))
{
cmd.ExecuteNonQuery();
}
Console.WriteLine($"已向表 '{tableName}' 添加列 '{columnName}'");
}
// 检查表是否存在
private bool TableExists(string tableName)
{
string query = "SELECT count(*) FROM information_schema.tables WHERE table_schema = 'public' AND table_name = @tableName;";
using (var cmd = new NpgsqlCommand(query, db))
{
cmd.Parameters.AddWithValue("@tableName", tableName);
var result = cmd.ExecuteScalar();
return Convert.ToInt64(result) > 0;
}
}
// 检查列是否已存在
private bool ColumnExists(string tableName, string columnName)
{
string query = "SELECT column_name FROM information_schema.columns WHERE table_schema = 'public' AND table_name = @tableName;";
using (var cmd = new NpgsqlCommand(query, db))
{
cmd.Parameters.AddWithValue("@tableName", tableName);
using (var reader = cmd.ExecuteReader())
{
while (reader.Read())
{
if (reader["column_name"].ToString().Equals(columnName, StringComparison.OrdinalIgnoreCase))
{
return true;
}
}
}
}
return false;
}
// 批量添加多个列
public void AddColumns(string tableName, Dictionary<string, string> columns)
{
foreach (var column in columns)
{
AddColumn(tableName, column.Key, column.Value);
}
}
// 插入数据:向指定表插入一条记录
// 参数: tableName - 表名
// 参数: columnValues - 字段和对应值的字典
// 例如: InsertData("Users", new Dictionary<string, object> { {"Name", "John"}, {"Age", 30} });
public void InsertData(string tableName, Dictionary<string, object> columnValues)
{
// 构建列名字符串(添加引号防止保留关键字冲突)
var columns = string.Join(", ", columnValues.Keys.Select(k => $"\"{k}\""));
// 判断是否需要 JSON 转换,并构建参数占位符
var parameters = columnValues.Keys.Select(k =>
{
var value = columnValues[k];
if (value is JObject || value is JArray) // 判断是否是 JObject 或 JArray
{
return $"@{k}::json";
}
return $"@{k}";
});
// 构建插入语句
string insertQuery = $"INSERT INTO \"{tableName}\" ({columns}) VALUES ({string.Join(", ", parameters)})";
// 使用 NpgsqlCommand 而不是 SQLiteCommand
using (var cmd = new NpgsqlCommand(insertQuery, db))
{
foreach (var kvp in columnValues)
{
if (kvp.Value is JObject jObj)
{
// JObject 转为字符串
cmd.Parameters.AddWithValue("@" + kvp.Key, jObj.ToString());
}
else if (kvp.Value is JArray jArray)
{
// JArray 转为字符串
cmd.Parameters.AddWithValue("@" + kvp.Key, jArray.ToString());
}
else
{
// 处理 null 值
cmd.Parameters.AddWithValue("@" + kvp.Key, kvp.Value ?? DBNull.Value);
}
}
cmd.ExecuteNonQuery();
}
}
// 查询数据:执行任意查询语句并返回结果
public List<Dictionary<string, object>> SelectData(string query, Dictionary<string, object> parameters = null)
{
var result = new List<Dictionary<string, object>>();
using (var conn = new NpgsqlConnection(_connectionString))
using (var cmd = new NpgsqlCommand(query, db))
{
try
// 添加查询参数(如果有的话)
if (parameters != null)
{
await conn.OpenAsync();
Debug.WriteLine("Database connection established.");
using (var cmd = new NpgsqlCommand(query, conn))
foreach (var kvp in parameters)
{
if (parameters != null)
{
foreach (var param in parameters)
{
cmd.Parameters.AddWithValue(param.Key, param.Value ?? DBNull.Value);
}
}
using (var reader = await cmd.ExecuteReaderAsync())
{
while (await reader.ReadAsync())
{
var row = new Dictionary<string, object>();
for (int i = 0; i < reader.FieldCount; i++)
{
row[reader.GetName(i)] = await reader.IsDBNullAsync(i) ? null : reader.GetValue(i);
}
result.Add(row);
}
}
cmd.Parameters.AddWithValue("@" + kvp.Key, kvp.Value ?? DBNull.Value);
}
Debug.WriteLine($"Query executed: {query}");
}
catch (Exception ex)
using (var reader = cmd.ExecuteReader())
{
Debug.WriteLine($"Error executing query: {ex.Message}");
while (reader.Read())
{
var row = new Dictionary<string, object>();
for (int i = 0; i < reader.FieldCount; i++)
{
row[reader.GetName(i)] = reader.GetValue(i);
}
result.Add(row);
}
}
}
return result;
}
// 执行插入、更新、删除操作
public async Task<int> ExecuteNonQueryAsync(string query, Dictionary<string, object> parameters = null)
// 更新数据:根据条件更新指定表中的记录
public void UpdateData(string tableName, Dictionary<string, object> columnValues, string condition)
{
using (var conn = new NpgsqlConnection(_connectionString))
// 构建SET子句
var setClause = string.Join(", ", columnValues.Keys.Select(k => $"\"{k}\" = @{k}"));
string updateQuery = $"UPDATE \"{tableName}\" SET {setClause} WHERE {condition}";
using (var cmd = new NpgsqlCommand(updateQuery, db))
{
foreach (var kvp in columnValues)
{
cmd.Parameters.AddWithValue("@" + kvp.Key, kvp.Value ?? DBNull.Value);
}
cmd.ExecuteNonQuery();
}
}
// 删除数据:根据条件删除指定表中的记录
public void DeleteData(string tableName, string condition)
{
string deleteQuery = $"DELETE FROM \"{tableName}\" WHERE {condition}";
using (var cmd = new NpgsqlCommand(deleteQuery, db))
{
cmd.ExecuteNonQuery();
}
}
// 支持事务操作:允许在同一个事务中执行多个操作
public void ExecuteTransaction(Action transactionActions)
{
using (var transaction = db.BeginTransaction())
{
try
{
await conn.OpenAsync();
Debug.WriteLine("Database connection established.");
using (var cmd = new NpgsqlCommand(query, conn))
{
if (parameters != null)
{
foreach (var param in parameters)
{
cmd.Parameters.AddWithValue(param.Key, param.Value ?? DBNull.Value);
}
}
var rowsAffected = await cmd.ExecuteNonQueryAsync();
Debug.WriteLine($"Executed query: {query}, Rows affected: {rowsAffected}");
return rowsAffected;
}
transactionActions.Invoke(); // 执行多个操作
transaction.Commit(); // 提交事务
}
catch (Exception ex)
catch (Exception)
{
Debug.WriteLine($"Error executing non-query: {ex.Message}");
return -1;
transaction.Rollback(); // 回滚事务
throw;
}
}
}
// 执行插入操作,返回生成的主键
public async Task<int> ExecuteInsertAsync(string query, Dictionary<string, object> parameters = null, string returnColumn = "id")
// 删除所有表
public void DropAllTables()
{
using (var conn = new NpgsqlConnection(_connectionString))
// 获取所有表名
string getTablesQuery = "SELECT table_name FROM information_schema.tables WHERE table_schema = 'public';";
var tables = SelectData(getTablesQuery);
foreach (var table in tables)
{
try
string tableName = table["table_name"].ToString();
string dropTableQuery = $"DROP TABLE IF EXISTS \"{tableName}\" CASCADE;";
using (var cmd = new NpgsqlCommand(dropTableQuery, db))
{
await conn.OpenAsync();
Debug.WriteLine("Database connection established.");
using (var cmd = new NpgsqlCommand(query, conn))
{
if (parameters != null)
{
foreach (var param in parameters)
{
cmd.Parameters.AddWithValue(param.Key, param.Value ?? DBNull.Value);
}
}
cmd.CommandText += $" RETURNING {returnColumn};";
var result = await cmd.ExecuteScalarAsync();
Debug.WriteLine($"Executed insert, inserted ID: {result}");
return result != null ? Convert.ToInt32(result) : -1;
}
}
catch (Exception ex)
{
Debug.WriteLine($"Error executing insert: {ex.Message}");
return -1;
cmd.ExecuteNonQuery();
}
}
}
// 执行事务操作
public async Task<bool> ExecuteTransactionAsync(List<string> queries, List<Dictionary<string, object>> parametersList)
public int ExecuteBatchInsert(string tableName, List<string> columns, string columnsStr, List<Dictionary<string, object>> batch, int batchIndex)
{
using (var conn = new NpgsqlConnection(_connectionString))
var valueParams = new List<string>();
var parameters = new Dictionary<string, object>();
for (int i = 0; i < batch.Count; i++)
{
try
{
await conn.OpenAsync();
Debug.WriteLine("Database connection established.");
using (var transaction = await conn.BeginTransactionAsync())
{
for (int i = 0; i < queries.Count; i++)
{
using (var cmd = new NpgsqlCommand(queries[i], conn, transaction))
{
var parameters = parametersList[i];
if (parameters != null)
{
foreach (var param in parameters)
{
cmd.Parameters.AddWithValue(param.Key, param.Value ?? DBNull.Value);
}
}
var paramNames = columns.Select(col => $"@p{batchIndex}_{i}_{col}").ToList();
valueParams.Add($"({string.Join(", ", paramNames)})");
await cmd.ExecuteNonQueryAsync();
}
}
await transaction.CommitAsync();
Debug.WriteLine("Transaction committed.");
return true;
}
}
catch (Exception ex)
foreach (var col in columns)
{
Debug.WriteLine($"Error executing transaction: {ex.Message}");
return false;
parameters[$"p{batchIndex}_{i}_{col}"] = batch[i][col] ?? DBNull.Value;
}
}
string insertQuery = $"INSERT INTO \"{tableName}\" ({columnsStr}) VALUES {string.Join(", ", valueParams)}";
using (var transaction = db.BeginTransaction())
using (var cmd = new NpgsqlCommand(insertQuery, db, transaction))
{
foreach (var param in parameters)
{
cmd.Parameters.AddWithValue(param.Key, param.Value);
}
int result = cmd.ExecuteNonQuery();
transaction.Commit();
return result;
}
}
// 修改主方法
public int BulkInsert(string tableName, List<Dictionary<string, object>> dataList)
{
if (dataList == null || dataList.Count == 0)
return 0;
int insertedCount = 0;
var columns = dataList[0].Keys.ToList();
var columnsStr = string.Join(", ", columns.Select(c => $"\"{c}\""));
int batchSize = 500;
int batchIndex = 0;
for (int i = 0; i < dataList.Count; i += batchSize)
{
var batch = dataList.Skip(i).Take(batchSize).ToList();
insertedCount += ExecuteBatchInsert(tableName, columns, columnsStr, batch, batchIndex);
batchIndex++;
}
return insertedCount;
}
// 释放资源,关闭数据库连接
public void Dispose()
{
if (!disposed)
{
if (db != null && db.State == ConnectionState.Open)
{
db.Close();
db.Dispose();
}
disposed = true;
}
GC.SuppressFinalize(this);
}
// 析构函数调用Dispose释放资源
~PostgreSQL()
{
Dispose();
}
// PostgreSQL特有的方法创建连接字符串的便捷方法
public static string BuildConnectionString(string host, string database, string username, string password, int port = 5432)
{
return $"Host={host};Port={port};Database={database};Username={username};Password={password}";
}
}
}
}

View File

@@ -1,4 +1,5 @@
using System.Data;
using Newtonsoft.Json.Linq;
using System.Data;
using System.Data.SQLite;
using static NPOI.HSSF.Util.HSSFColor;
@@ -263,6 +264,60 @@ namespace Ramitta
}
}
public int ExecuteBatchInsert(string tableName, List<string> columns, string columnsStr, List<Dictionary<string, object>> batch, int batchIndex)
{
var valueParams = new List<string>();
var parameters = new Dictionary<string, object>();
for (int i = 0; i < batch.Count; i++)
{
var paramNames = columns.Select(col => $"@p{batchIndex}_{i}_{col}").ToList();
valueParams.Add($"({string.Join(", ", paramNames)})");
foreach (var col in columns)
{
parameters[$"p{batchIndex}_{i}_{col}"] = batch[i][col] ?? DBNull.Value;
}
}
string insertQuery = $"INSERT INTO {tableName} ({columnsStr}) VALUES {string.Join(", ", valueParams)}";
using (var transaction = db.BeginTransaction())
using (var cmd = new SQLiteCommand(insertQuery, db, transaction))
{
foreach (var param in parameters)
{
cmd.Parameters.AddWithValue(param.Key, param.Value);
}
int result = cmd.ExecuteNonQuery();
transaction.Commit();
return result;
}
}
// 修改主方法
public int BulkInsert(string tableName, List<Dictionary<string, object>> dataList)
{
if (dataList == null || dataList.Count == 0)
return 0;
int insertedCount = 0;
var columns = dataList[0].Keys.ToList();
var columnsStr = string.Join(", ", columns);
int batchSize = 500;
int batchIndex = 0; // 添加批次索引
for (int i = 0; i < dataList.Count; i += batchSize)
{
var batch = dataList.Skip(i).Take(batchSize).ToList();
insertedCount += ExecuteBatchInsert(tableName, columns, columnsStr, batch, batchIndex);
batchIndex++; // 递增批次索引
}
return insertedCount;
}
// 释放资源,关闭数据库连接
// 确保数据库连接在对象销毁时被正确关闭
public void Dispose()