///
/// 同步操作数据类
///
using MySql.Data.MySqlClient;
using System;
using System.Collections.Generic;
using System.Data;
using System.Data.SqlClient;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
using static AFKAbutment.EnumClass;
namespace AFKAbutment
{
public class SynchData
{
public string mysqlConnectionString;//mysql连接字段
public string sqlServerConnectionString;//sqlsever连接字段
///
/// 说明:创建Afk向朗速传递数据的同步方法
/// 创建人:王一帆
/// 创建日期:2023-11-10
/// 修改人:
/// 修改日期:
/// 修改备注:
/// 版本:1.0
///
/// Ls方需要构建的表
/// Afk方需要构建的表
/// 特殊不同字段集合
public void CreateAndSyncTable(string sqlServerTableName, string mysqlTableName, List specialFieldNames=null,DataTable LstoAfkData=null)
{
using (var sourceConnection = new MySqlConnection(mysqlConnectionString))
{
sourceConnection.Open();
if (LstoAfkData == null)
{
// 获取MySQL表结构
DataTable schemaTable = GetSchemaTable(mysqlTableName, sourceConnection);
// 创建SQL Server表并添加特殊字段
using (var destinationConnection = new SqlConnection(sqlServerConnectionString))
{
destinationConnection.Open();
if (!CheckIfTableExists(sqlServerTableName, destinationConnection))
{
CreateSqlServerTableWithSpecialFields(sqlServerTableName, schemaTable, specialFieldNames, destinationConnection);
}
// 同步数据
using (var bulkCopy = new SqlBulkCopy(destinationConnection))
{
bulkCopy.DestinationTableName = sqlServerTableName;
//sourceConnection.ChangeDatabase("mysqlDatabaseName"); // 替换为MySQL数据库名
using (var selectCmd = new MySqlCommand($"SELECT * FROM {mysqlTableName}", sourceConnection))
using (var reader = selectCmd.ExecuteReader())
{
bulkCopy.WriteToServer(reader);
}
}
}
}
else//若存在数据源则反向传数据源
{
}
}
}
///
/// 说明:获取mysql表中的列信息方便调用时候直接取用
/// 创建人:王一帆
/// 创建日期:2023-11-10
/// 修改人:
/// 修改日期:
/// 修改备注:
/// 版本:1.0
///
/// mysql表名
/// MySqlConnection的连接
///
private DataTable GetSchemaTable(string tableName, MySqlConnection connection)
{
using (var schemaCommand = new MySqlCommand($"DESCRIBE {tableName}", connection))
{
using (var reader = schemaCommand.ExecuteReader())
{
DataTable schemaTable = new DataTable();
schemaTable.Load(reader);
return schemaTable;
}
}
}
///
/// 说明:判断Ls数据库中是否存在需要同步的表名
/// 创建人:王一帆
/// 创建日期:2023-11-10
/// 修改人:
/// 修改日期:
/// 修改备注:
/// 版本:1.0
///
/// 表名
/// SqlConnection的连接
///
private bool CheckIfTableExists(string tableName, SqlConnection connection)
{
using (var command = new SqlCommand($"SELECT CASE WHEN EXISTS((SELECT * FROM INFORMATION_SCHEMA.TABLES WHERE TABLE_NAME = '{tableName}')) THEN 1 ELSE 0 END", connection))
{
return (int)command.ExecuteScalar() == 1;
}
}
///
/// 说明:通过对应的数据结构进行同步
/// 创建人:王一帆
/// 创建日期:2023-11-10
/// 修改人:
/// 修改日期:
/// 修改备注:
/// 版本:1.0
///
/// Ls方需要构建的表名
///
///
///
private void CreateSqlServerTableWithSpecialFields(string tableName, DataTable schemaTable, List specialFieldNames, SqlConnection connection)
{
using (var command = new SqlCommand("", connection))
{
// 创建表的SQL语句
command.CommandText = $"CREATE TABLE {tableName} (";
foreach (DataRow row in schemaTable.Rows)
{
var columnName = row["Field"].ToString();
var columnType = row["Type"].ToString();
// 转换MySQL数据类型到SQL Server数据类型
if (columnType.StartsWith("varchar"))
{
// 如果字段在特殊字段列表中或类型为varchar,则设置为NVARCHAR(255)(或其他需要的长度)
command.CommandText += $"[{columnName}] NVARCHAR(255), ";
}
else if (columnType.StartsWith("int"))
{
command.CommandText += $"[{columnName}] INT, ";
}
else if (columnType.StartsWith("decimal"))
{
command.CommandText += $"[{columnName}] DECIMAL(18, 2), "; // 适当调整精度和小数位
}
// 其他类型转换...
else
{
// 默认转换为NVARCHAR(MAX)
command.CommandText += $"[{columnName}] NVARCHAR(MAX), ";
}
}
command.CommandText = command.CommandText.TrimEnd(',', ' ') + ")";
command.ExecuteNonQuery();
}
}
}
}