ab56a9bcf7
SVN-Revision: r240
171 lines
7.2 KiB
C#
171 lines
7.2 KiB
C#
/// <summary>
|
|
/// 同步操作数据类
|
|
/// </summary>
|
|
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连接字段
|
|
/// <summary>
|
|
/// <para>说明:创建Afk向朗速传递数据的同步方法</para>
|
|
/// <para>创建人:王一帆</para>
|
|
/// <para>创建日期:2023-11-10 </para>
|
|
/// <para>修改人:</para>
|
|
/// <para>修改日期:</para>
|
|
/// <para>修改备注:</para>
|
|
/// <para>版本:1.0</para>
|
|
/// </summary>
|
|
/// <param name="sqlServerTableName">Ls方需要构建的表</param>
|
|
/// <param name="mysqlTableName">Afk方需要构建的表</param>
|
|
/// <param name="specialFieldNames">特殊不同字段集合</param>
|
|
public void CreateAndSyncTable(string sqlServerTableName, string mysqlTableName, List<string> 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//若存在数据源则反向传数据源
|
|
{
|
|
|
|
|
|
|
|
}
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// <para>说明:获取mysql表中的列信息方便调用时候直接取用</para>
|
|
/// <para>创建人:王一帆</para>
|
|
/// <para>创建日期:2023-11-10 </para>
|
|
/// <para>修改人:</para>
|
|
/// <para>修改日期:</para>
|
|
/// <para>修改备注:</para>
|
|
/// <para>版本:1.0</para>
|
|
/// </summary>
|
|
/// <param name="tableName">mysql表名</param>
|
|
/// <param name="connection">MySqlConnection的连接</param>
|
|
/// <returns></returns>
|
|
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;
|
|
}
|
|
}
|
|
}
|
|
/// <summary>
|
|
/// <para>说明:判断Ls数据库中是否存在需要同步的表名</para>
|
|
/// <para>创建人:王一帆</para>
|
|
/// <para>创建日期:2023-11-10 </para>
|
|
/// <para>修改人:</para>
|
|
/// <para>修改日期:</para>
|
|
/// <para>修改备注:</para>
|
|
/// <para>版本:1.0</para>
|
|
/// </summary>
|
|
/// <param name="tableName">表名</param>
|
|
/// <param name="connection">SqlConnection的连接</param>
|
|
/// <returns></returns>
|
|
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;
|
|
}
|
|
}
|
|
/// <summary>
|
|
/// <para>说明:通过对应的数据结构进行同步</para>
|
|
/// <para>创建人:王一帆</para>
|
|
/// <para>创建日期:2023-11-10 </para>
|
|
/// <para>修改人:</para>
|
|
/// <para>修改日期:</para>
|
|
/// <para>修改备注:</para>
|
|
/// <para>版本:1.0</para>
|
|
/// </summary>
|
|
/// <param name="tableName">Ls方需要构建的表名</param>
|
|
/// <param name="schemaTable"></param>
|
|
/// <param name="specialFieldNames"></param>
|
|
/// <param name="connection"></param>
|
|
private void CreateSqlServerTableWithSpecialFields(string tableName, DataTable schemaTable, List<string> 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();
|
|
}
|
|
}
|
|
|
|
}
|
|
}
|