C#中将数据批量导入数据库

本文讨论如何将DataTable中的数据批量导入数据库,并实现相关组件的封装。

tBatch接口

下面的代码(cfx/data/tBatch.cs)定义了tBatch接口类型,作为批量导入数据操作的标准。

C#
using System.Collections.Generic;
using System.Data;

namespace cfx.data
{
    public interface tBatch
    {
        string CnnStr { get; }
        string Table { get; }
        DataTable Data { get; }
        tBatch SetPrimaryKey(params string[] field);
        //
        long Insert();
        long TruncateAndInsert();
        long Update();
        long InsertOrUpdate();
        long ClearAndUpdate();
    }
}

tBatch接口中,首先定义了三个只读属性,分别是数据库连接字符串(CnnStr)、数据表名(Table)、包含导入数据的DataTable对象(Data)。接下来的SetPrimaryKey()方法用于指定查询或更新记录时的主键字段,可以设置一个或多个主键字段。最后是五个导入方法,分别是:

  • Insert(),向表中添加记录。
  • TruncateAndInsert(),重置表后添加记录。此操作会清空数据表后再添加新的数据。
  • Update(),根据主键字段更新记录。
  • InsertOrUpdate(),根据主键字段判断记录是否存在,存在时更新数据,不存在时添加数据。
  • ClearAndUpdate(),首先将表(Table)中Data对象定义的非主键字段设置为空(null),然后根据主键数据更新记录。

这五个方法都会返回最后添加或修改的记录数量,返回类型为long。

下面的代码(cfx/data/tBatch.cs)定义了tBatchBase类,是实现tBatch组件的基类。

C#
using System.Collections.Generic;
using System.Data;

namespace cfx.data
{
    // 其它代码
    public abstract class tBatchBase : tBatch
    {
        //
        public tBatchBase(string cnnstr,string table,DataTable data)
        {
            CnnStr = cnnstr;
            Table = table;
            Data = data;
        }
        //
        protected List<string> myPk;
        //
        public string CnnStr { get; private set; }
        public string Table { get; private set; }
        public DataTable Data { get; private set; }
        //
        public tBatch SetPrimaryKey(params string[] field)
        {
            if (field.Length == 0 || Data == null || Data.Columns.Count == 0)
            {
                myPk = null;
            }
            else
            {
                myPk = new List<string>(field);
                // 检查主键字段并移动到列的最后
                int colLastIndex = Data.Columns.Count - 1;
                for (int i = 0; i < myPk.Count; i++)
                {
                    if (Data.Columns.Contains(myPk[i]) == false)
                    {
                        // 主键列不存在则主键设置失败
                        myPk = null;
                        break;
                    }
                    //
                    Data.Columns[myPk[i]].SetOrdinal(colLastIndex);
                }
            }
            return this;
        }
        //
        public abstract long Insert();
        public abstract long TruncateAndInsert();
        public abstract long Update();
        public abstract long InsertOrUpdate();
        public abstract long ClearAndUpdate();
    }
    //
}

这里要注意SetPrimaryKey()方法的实现,方法中会在Data对象的列中检查指定的主键列是否存在,如果不存在则设置主键失败,myPk字段设置为null;如果主键列存在,则全部移动到最后,Data对象中主键列的顺序与myPk中的顺序相同。

tMySqlBatch类

下面的代码(cfx/data/mysql/tMySqlBatch.cs)实现了MySQL数据库批量导入数据的tMySqlBatch类。

C#
using System.Text;
using System.Data;
using MySql.Data.MySqlClient;

namespace cfx.data.mysql
{
    public class tMySqlBatch :tBatchBase
    {
        public tMySqlBatch(string cnnstr, string table, DataTable data)
            : base(cnnstr, table, data) { }
        //
        public static tMySqlBatch Create(string cnnstr,string table,DataTable data)
        {
            return new tMySqlBatch(cnnstr, table, data);
        }
        // 辅助方法
        // 添加记录的insert语句
        private string GetInsertSql()
        {
            if (CnnStr == null || CnnStr.Length == 0 ||
                Table == null || Table.Length == 0 ||
                Data == null || Data.Columns.Count == 0) return "";
            //
            StringBuilder sb = new StringBuilder(512);
            StringBuilder sbVal = new StringBuilder(256);
            sb.AppendFormat("insert into `{0}`(`{1}`",
                Table, Data.Columns[0].ColumnName);
            sbVal.Append(")values(?data0");
            for(int col = 1; col < Data.Columns.Count; col++)
            {
                sb.AppendFormat(",`{0}`", Data.Columns[col].ColumnName);
                sbVal.AppendFormat(",?data{0}", col);
            }
            sb.Append(sbVal.ToString());
            sb.Append(")");
            return sb.ToString();
        }
        // 根据主键更新记录的update语句
        private string GetUpdateSql()
        {
            if (CnnStr == null || CnnStr.Length == 0 ||
                Table == null || Table.Length == 0 ||
                myPk == null) return "";
            //
            int dataCount = Data.Columns.Count - myPk.Count;
            StringBuilder sb = new StringBuilder(256);
            sb.AppendFormat("update `{0}` set `{1}`=?data0",
                Table, Data.Columns[0].ColumnName);
            for(int col = 1; col < dataCount; col++)
            {
                sb.AppendFormat(",`{0}`=?data{1}",
                    Data.Columns[col].ColumnName, col);
            }
            //
            sb.AppendFormat(" where `{0}`=?data{1}",
                Data.Columns[dataCount].ColumnName, dataCount);
            for(int col = dataCount + 1; col < Data.Columns.Count; col++)
            {
                sb.AppendFormat(" and `{0}`=?data{1}",
                    Data.Columns[col].ColumnName, col);
            }
            //
            return sb.ToString();
        }
        // 判断主键数据记录是否存在的select语句
        private string GetExistsSql()
        {
            int dataCount = Data.Columns.Count - myPk.Count;
            StringBuilder sb = new StringBuilder(256);
            sb.AppendFormat("select exists(select * from `{0}`", Table);
            //
            sb.AppendFormat(" where `{0}`=?data{1}",
                Data.Columns[dataCount].ColumnName, dataCount);
            for (int col = dataCount + 1; col < Data.Columns.Count; col++)
            {
                sb.AppendFormat(" and `{0}`=?data{1}",
                    Data.Columns[col].ColumnName, col);
            }
            sb.Append(")");
            //
            return sb.ToString();
        }
        // 清除表中所有记录非主键(myPk)字段的数据,update语句
        private string GetClearSql()
        {
            int dataCount = Data.Columns.Count - myPk.Count;
            StringBuilder sb = new StringBuilder(256);
            sb.AppendFormat("update `{0}` set `{1}`=null",
                Table, Data.Columns[0].ColumnName);
            for(int col = 1; col < dataCount; col++)
            {
                sb.AppendFormat(",`{0}`=null", 
                    Data.Columns[col].ColumnName);
            }
            //
            return sb.ToString();
        }
        //
        // 添加记录
        public override long Insert()
        {
            try
            {
                string insSql = GetInsertSql();
                if (insSql.Length == 0) return -1001;
                using(MySqlConnection cnn = new MySqlConnection(CnnStr))
                {
                    cnn.Open();
                    MySqlCommand cmd = cnn.CreateCommand();
                    cmd.CommandText = insSql;
                    using (MySqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        long counter = 0;
                        for(int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            for(int col = 0; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("?data" + col,
                                    Data.Rows[row][col]);
                            }
                            counter += cmd.ExecuteNonQueryAsync().Result;
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch
            {
                return -1000;
            }
        }
        // 清表后添加记录
        public override long TruncateAndInsert()
        {
            try
            {
                string insSql = GetInsertSql();
                if (insSql.Length == 0) return -1001;
                //
                using (MySqlConnection cnn = new MySqlConnection(CnnStr))
                {
                    cnn.Open();
                    MySqlCommand cmd = cnn.CreateCommand();
                    using (MySqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        // 清表
                        cmd.CommandText = "truncate `"+Table+"`";
                        cmd.ExecuteNonQuery();
                        //
                        cmd.CommandText = insSql;
                        long counter = 0;
                        for (int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            for (int col = 0; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("?data" + col,
                                    Data.Rows[row][col]);
                            }
                            counter += cmd.ExecuteNonQuery();
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch
            {
                return -1000;
            }
        }
        // 根据主键更新记录
        public override long Update()
        {
            try
            {
                string updSql = GetUpdateSql();
                if (updSql.Length == 0) return -1001;
                //
                using (MySqlConnection cnn = new MySqlConnection(CnnStr))
                {
                    cnn.Open();
                    MySqlCommand cmd = cnn.CreateCommand();
                    cmd.CommandText = updSql;
                    using (MySqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        long counter = 0;
                        for (int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            for (int col = 0; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("?data" + col,
                                    Data.Rows[row][col]);
                            }
                            counter += cmd.ExecuteNonQuery();
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch
            {
                return -1000;
            }
        }
    
        // 根据主键判断,记录存在则更新,不存在则添加
        public override long InsertOrUpdate()
        {
            try
            {
                string insSql = GetInsertSql();
                string updSql = GetUpdateSql();
                string selSql = GetExistsSql();
                if (insSql.Length == 0 || updSql.Length == 0)
                    return -1001;
                //
                int dataCount = Data.Columns.Count - myPk.Count;
                using (MySqlConnection cnn = new MySqlConnection(CnnStr))
                {
                    cnn.Open();
                    MySqlCommand cmd = cnn.CreateCommand();
                    using (MySqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        long counter = 0;
                        for (int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            // 判断记录是否存在
                            cmd.CommandText = selSql;
                            for(int col = dataCount; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("?data" + col,
                                    Data.Rows[row][col]);
                            }
                            if (Obj.ToInt(cmd.ExecuteScalar()) == 1)
                            {
                                // 存在,更新
                                cmd.Parameters.Clear();
                                cmd.CommandText = updSql;
                                for (int col = 0; col < Data.Columns.Count; col++)
                                {
                                    cmd.Parameters.AddWithValue("?data" + col,
                                        Data.Rows[row][col]);
                                }
                                counter += cmd.ExecuteNonQuery();
                            }
                            else
                            {
                                // 不存在,添加
                                cmd.Parameters.Clear();
                                cmd.CommandText = insSql;
                                for (int col = 0; col < Data.Columns.Count; col++)
                                {
                                    cmd.Parameters.AddWithValue("?data" + col,
                                        Data.Rows[row][col]);
                                }
                                counter += cmd.ExecuteNonQuery();
                            }
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch { return -1000; }
        }
        // 清除非主键字段数据后根据主键数据更新记录
        public override long ClearAndUpdate()
        {
            try
            {
                string updSql = GetUpdateSql();
                string clearSql = GetClearSql();
                if (updSql.Length == 0) return -1001;
                //
                using(MySqlConnection cnn = new MySqlConnection(CnnStr))
                {
                    cnn.Open();
                    MySqlCommand cmd = cnn.CreateCommand();
                    using (MySqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        // 清理数据
                        cmd.CommandText = clearSql;
                        cmd.ExecuteNonQuery();
                        // 更新数据
                        cmd.CommandText = updSql;
                        long counter = 0;
                        for(int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            for (int col = 0; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("?data" + col,
                                    Data.Rows[row][col]);
                            }
                            counter += cmd.ExecuteNonQuery();
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch { return -1000; }
        }
        //
        //
        //
    }
}

tMySqlBatch类中,构造函数和Create()静态方法用于创建tMySqlBatch类的实例,都需要三个参数,分别是数据库连接字符串、数据表名称、包含导入数据的DataTable对象。

下面分别解释各种导入操作方法的实现。

添加数据

Insert()和TruncateAndInsert()方法都用于在表中添加数据记录,不同的是,TruncateAndInsert()方法会使用truncate语句重置表,此操作会清除表中的所有数据,记录ID也会重新开始计数。

添加数据的辅助方法为GetInsertSql(),用于生成insert语句,方法中首先对CnnStr、Table和Data进行了检查,如果不满足操作要求则返回空字符串。方法中,字段名使用了Data对象中的列名,参数名称使用?data0、?data1、?data2、……。方法返回的insert语句格式如下:

MySQL
insert into `表名`(`列名0`,`列名1`,`列名2`,...)values(?data0,?data1,?data2,...)

Insert()和TruncateAndInsert()方法中,通过事务执行添加操作,Data对象中的每一行执行一次添加操作,添加的记录数量会使用counter变量计数。最后,方法会返回counter变量的值,即添加了多少条记录。方法中,如果操作异常返回-1000,操作参数不满足要求则返回-1001。

更新数据

更新数据的语句使用GetUpdateSql()方法创建,其中,对数据库连接字符串、数据表和主键(myPk)进行了检查,不满足操作要求时返回空字符串。使用SetPrimaryKey()方法设置主键时,修改了Data对象中列的位置,主键列都排到后面,在组件更新条件和添加语句的参数数据时会比较方便。GetUpdateSql()方法中,需要注意dataCount变量,其值是Data对象不包含主键的数据列数量,整理后的Data对象中,数据列的索引从0到dataCount-1;主键列索引从dataCount到Data.Columns.Count-1,即update语句中使用的条件数据列。

Update()方法中,同样会先获取操作语句,当其为空字符串时,也就是不满足操作要求时,方法返回-1001。接下来,会在事务中逐条更新Data中的数据,并将每次操作影响的记录数量保存到counter变量中,操作成功后会返回counter变量的值,操作异常时同样返回-1000。

添加或更新

InsertOrUpdate()方法会根据myPk中指定的主键数据判断数据表中是否存在指定的记录,如果记录已存在则更新为Data对象中的数据,如果记录不存在,则添加Data对象中的数据。

判断记录是否存在时使用了MySQL中的exists()函数,其select语句由GetExistsSql()方法生成,语句会返回exists()函数的返回值,当exists()函数中的select语句包含查询结果时返回1,否则返回0。

InsertOrUpdate()方法中会使用三条语句,分别是添加记录的insert语句、修改记录的update语句、判断记录是否存在的select语句。

处理Data对象中的行记录数据时,首先会通过select语句判断主键数据的记录是否存在,当记录存在时执行update语句更新数据,不存在时执行insert语句添加记录。方法中,counter变量保存的是更新或添加的全部记录数量。

此操作中,如果需要分别获取更新和添加的记录数据,可以通过两个计数变量分别保存,然后通过输出参数输出,这里需要修改tBatch接口的定义及相关类的实现。

清除并更新

ClearAndUpdate()方法的功能是:首先将Data对象中非主键列在表中对应字段的数据设置为null,然后根据Data中的数据逐条记录更新。方法中使用了update语句,包括清除原数据和逐行更新Data对象数据的操作,其中,清除数据的语句由GetClearSql()方法创建。方法会返回更新的记录数量,执行异常时返回-1000。

测试tMySqlBatch类

接下来,测试tMySqlBatch类时会继续使用cdb_cs1数据库中的t1表,相关定义可以参考:http://caohuayu.com/article/Article.aspx?id=a264020。

首先,保持t1表的原始状态,如果t1表中已存在数据,可以在HeidiSQL中执行如下代码重置。

MySQL
use cdb_cs1;
truncate t1;

然后在d:\t1.xls文件中准备需要导入的数据,如下图所示。(示例中使用的Excel文件可以从附件中下载)。

准备导入数据

这里只需要准备f1、f2、f3、f4字段数据,recid为自动ID,可以自动添加。接下来,修改Program.cs文件代码如下。

C#
using System;
using System.Data;
using cfx.office;
using cfx.data;
using cfx.data.mysql;

namespace csfx_demo
{
    class Program
    {
        static void Main(string[] args)
        {
            DataTable data = tExcelReader.Create(@"d:\t1.xls").Read();
            string cnnstr =
    tMySqlHelper.GetCnnStr("127.0.0.1", "cdb_cs1", "root", "DEV_Test123456", 3306);
            tBatch bat = tMySqlBatch.Create(cnnstr, "t1", data);
            long result = bat.Insert();
            Console.WriteLine(result);
        }
    }
}

代码会读取d:\t1.xls文件第一个工作表的数据到data对象,然后,通过tMySqlBatch对象的Insert()方法导入到t1表,执行成功时会返回导入的记录数量10。读取Excel数据操作可以参考前一篇文章。

需要清除数据表中的所有数据,然后导入新的数据时可以使用TruncateAndInsert()方法,如下面的代码。

C#
using System;
using System.Data;
using cfx.office;
using cfx.data;
using cfx.data.mysql;

namespace csfx_demo
{
    class Program
    {
        static void Main(string[] args)
        {
            DataTable data = tExcelReader.Create(@"d:\t1.xls").Read();
            string cnnstr =
                tMySqlHelper.GetCnnStr("127.0.0.1", "cdb_cs1", "root", "DEV_Test123456", 3306);
            tBatch bat = tMySqlBatch.Create(cnnstr, "t1", data);
            long result = bat.TruncateAndInsert();
            Console.WriteLine(result);
        }
    }
}

下面在d:\t2.xls文件中准备如下图所示的数据。

准备更新数据

然后,修改Program.cs文件代码如下。

C#
using System;
using System.Data;
using cfx.office;
using cfx.data;
using cfx.data.mysql;

namespace csfx_demo
{
    class Program
    {
        static void Main(string[] args)
        {
            DataTable data = tExcelReader.Create(@"d:\t2.xls").Read();
            string cnnstr =
    tMySqlHelper.GetCnnStr("127.0.0.1", "cdb_cs1", "root", "DEV_Test123456", 3306);
            tBatch bat = tMySqlBatch.Create(cnnstr, "t1", data).SetPrimaryKey("f1");
            long result = bat.Update();
            Console.WriteLine(result);
        }
    }
}

代码中,首先读取d:\t2.xls文件的数据到data对象;然后会以f1字段数据为主键更新t1表,这里会修改f4字段的数据,修改后的t1表数据如下图所示。

批量更新数据

在d:\t3.xls文件中准备如下图所示的数据,这里新增了f1字段等于user11和user12的记录。

准备数据

修改Program.cs文件代码如下。

C#
using System;
using System.Data;
using cfx.office;
using cfx.data;
using cfx.data.mysql;

namespace csfx_demo
{
    class Program
    {
        static void Main(string[] args)
        {
            DataTable data = tExcelReader.Create(@"d:\t3.xls").Read();
            string cnnstr =
    tMySqlHelper.GetCnnStr("127.0.0.1", "cdb_cs1", "root", "DEV_Test123456", 3306);
            tBatch bat = tMySqlBatch.Create(cnnstr, "t1", data).SetPrimaryKey("f1");
            long result = bat.InsertOrUpdate();
            Console.WriteLine(result);
        }
    }
}

代码会读取d:\t3.xls文件的数据,然后使用f1字段为主键数据,当f1数据存在时则更新数据,否则添加新的记录,执行后t1表的数据如下图所示。

添加或更新数据

在d:\t4.xls文件中准备如下图所示的数据。

准备数据

修改Program.cs文件代码如下。

C#
using System;
using System.Data;
using cfx.office;
using cfx.data;
using cfx.data.mysql;

namespace csfx_demo
{
    class Program
    {
        static void Main(string[] args)
        {
            DataTable data = tExcelReader.Create(@"d:\t4.xls").Read();
            string cnnstr =
    tMySqlHelper.GetCnnStr("127.0.0.1", "cdb_cs1", "root", "DEV_Test123456", 3306);
            tBatch bat = tMySqlBatch.Create(cnnstr, "t1", data).SetPrimaryKey("f1");
            long result = bat.ClearAndUpdate();
            Console.WriteLine(result);
        }
    }
}

执行代码会清除t1表中f4字段的数据,然后更新d:\t4.xls文件中指定记录(f1为主键)的f4字段数据,操作完成后的t1表数据如下图所示。

清除字段并更新

可以看到,在ClearAndUpdate()方法操作后,由于d:\t4.xls文件中不包含f1字段为user06到user10的记录,所以,这些记录的f4字段数据清除后为空值。

使用ClearAndUpdate()方法时应注意,此方法更新的字段必须允许为空,在更新数据时应尽量设置少的字段数据。如果需要保留数据字段中原有的数据,应使用Update()方法。

tSqlBatch类

下面的代码(cfx/data/sql/tSqlBatch.cs)给出了tSqlBatch类的实现,用于SQL Server数据库的批量导入。

C#
using System.Text;
using System.Data;
using System.Data.SqlClient;

namespace cfx.data.sql
{
    public class tSqlBatch :tBatchBase
    {
        public tSqlBatch(string cnnstr, string table, DataTable data)
            : base(cnnstr, table, data) { }
        //
        public static tSqlBatch Create(string cnnstr,string table,DataTable data)
        {
            return new tSqlBatch(cnnstr, table, data);
        }
        // 辅助方法
        // 添加记录的insert语句
        private string GetInsertSql()
        {
            if (CnnStr == null || CnnStr.Length == 0 ||
                Table == null || Table.Length == 0 ||
                Data == null || Data.Columns.Count == 0) return "";
            //
            StringBuilder sb = new StringBuilder(512);
            StringBuilder sbVal = new StringBuilder(256);
            sb.AppendFormat("insert into [{0}]([{1}]",
                Table, Data.Columns[0].ColumnName);
            sbVal.Append(")values(@data0");
            for(int col = 1; col < Data.Columns.Count; col++)
            {
                sb.AppendFormat(",[{0}]", Data.Columns[col].ColumnName);
                sbVal.AppendFormat(",@data{0}", col);
            }
            sb.Append(sbVal.ToString());
            sb.Append(")");
            return sb.ToString();
        }
        // 根据主键更新记录的update语句
        private string GetUpdateSql()
        {
            if (CnnStr == null || CnnStr.Length == 0 ||
                Table == null || Table.Length == 0 ||
                myPk == null) return "";
            //
            int dataCount = Data.Columns.Count - myPk.Count;
            StringBuilder sb = new StringBuilder(256);
            sb.AppendFormat("update [{0}] set [{1}]=@data0",
                Table, Data.Columns[0].ColumnName);
            for(int col = 1; col < dataCount; col++)
            {
                sb.AppendFormat(",[{0}]=@data{1}",
                    Data.Columns[col].ColumnName, col);
            }
            //
            sb.AppendFormat(" where [{0}]=@data{1}",
                Data.Columns[dataCount].ColumnName, dataCount);
            for(int col = dataCount + 1; col < Data.Columns.Count; col++)
            {
                sb.AppendFormat(" and [{0}]=@data{1}",
                    Data.Columns[col].ColumnName, col);
            }
            //
            return sb.ToString();
        }
        // 判断主键数据记录是否存在的select语句
        private string GetExistsSql()
        {
            int dataCount = Data.Columns.Count - myPk.Count;
            StringBuilder sb = new StringBuilder(256);
            sb.AppendFormat("select 1 where exists(select * from [{0}]", Table);
            //
            sb.AppendFormat(" where [{0}]=@data{1}",
                Data.Columns[dataCount].ColumnName, dataCount);
            for (int col = dataCount + 1; col < Data.Columns.Count; col++)
            {
                sb.AppendFormat(" and [{0}]=@data{1}",
                    Data.Columns[col].ColumnName, col);
            }
            sb.Append(")");
            //
            return sb.ToString();
        }
        // 清除表中所有记录非主键(myPk)字段的数据,update语句
        private string GetClearSql()
        {
            int dataCount = Data.Columns.Count - myPk.Count;
            StringBuilder sb = new StringBuilder(256);
            sb.AppendFormat("update [{0}] set [{1}]=null",
                Table, Data.Columns[0].ColumnName);
            for(int col = 1; col < dataCount; col++)
            {
                sb.AppendFormat(",[{0}]=null", 
                    Data.Columns[col].ColumnName);
            }
            //
            return sb.ToString();
        }
        //
        // 添加记录
        public override long Insert()
        {
            try
            {
                string insSql = GetInsertSql();
                if (insSql.Length == 0) return -1001;
                using(SqlConnection cnn = new SqlConnection(CnnStr))
                {
                    cnn.Open();
                    SqlCommand cmd = cnn.CreateCommand();
                    cmd.CommandText = insSql;
                    using (SqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        long counter = 0;
                        for(int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            for(int col = 0; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("@data" + col,
                                    Data.Rows[row][col]);
                            }
                            counter += cmd.ExecuteNonQueryAsync().Result;
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch
            {
                return -1000;
            }
        }
        // 清表后添加记录
        public override long TruncateAndInsert()
        {
            try
            {
                string insSql = GetInsertSql();
                if (insSql.Length == 0) return -1001;
                //
                using (SqlConnection cnn = new SqlConnection(CnnStr))
                {
                    cnn.Open();
                    SqlCommand cmd = cnn.CreateCommand();
                    using (SqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        // 清表
                        cmd.CommandText = "truncate table ["+Table+"]";
                        cmd.ExecuteNonQuery();
                        //
                        cmd.CommandText = insSql;
                        long counter = 0;
                        for (int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            for (int col = 0; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("@data" + col,
                                    Data.Rows[row][col]);
                            }
                            counter += cmd.ExecuteNonQuery();
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch
            {
                return -1000;
            }
        }
        // 根据主键更新记录
        public override long Update()
        {
            try
            {
                string updSql = GetUpdateSql();
                if (updSql.Length == 0) return -1001;
                //
                using (SqlConnection cnn = new SqlConnection(CnnStr))
                {
                    cnn.Open();
                    SqlCommand cmd = cnn.CreateCommand();
                    cmd.CommandText = updSql;
                    using (SqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        long counter = 0;
                        for (int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            for (int col = 0; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("@data" + col,
                                    Data.Rows[row][col]);
                            }
                            counter += cmd.ExecuteNonQuery();
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch
            {
                return -1000;
            }
        }
    
        // 根据主键判断,记录存在则更新,不存在则添加
        public override long InsertOrUpdate()
        {
            try
            {
                string insSql = GetInsertSql();
                string updSql = GetUpdateSql();
                string selSql = GetExistsSql();
                if (insSql.Length == 0 || updSql.Length == 0)
                    return -1001;
                //
                int dataCount = Data.Columns.Count - myPk.Count;
                using (SqlConnection cnn = new SqlConnection(CnnStr))
                {
                    cnn.Open();
                    SqlCommand cmd = cnn.CreateCommand();
                    using (SqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        long counter = 0;
                        for (int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            // 判断记录是否存在
                            cmd.CommandText = selSql;
                            for(int col = dataCount; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("@data" + col,
                                    Data.Rows[row][col]);
                            }
                            if (Obj.ToInt(cmd.ExecuteScalar()) == 1)
                            {
                                // 存在,更新
                                cmd.Parameters.Clear();
                                cmd.CommandText = updSql;
                                for (int col = 0; col < Data.Columns.Count; col++)
                                {
                                    cmd.Parameters.AddWithValue("@data" + col,
                                        Data.Rows[row][col]);
                                }
                                counter += cmd.ExecuteNonQuery();
                            }
                            else
                            {
                                // 不存在,添加
                                cmd.Parameters.Clear();
                                cmd.CommandText = insSql;
                                for (int col = 0; col < Data.Columns.Count; col++)
                                {
                                    cmd.Parameters.AddWithValue("@data" + col,
                                        Data.Rows[row][col]);
                                }
                                counter += cmd.ExecuteNonQuery();
                            }
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch { return -1000; }
        }
        // 清除非主键字段数据后根据主键数据更新记录
        public override long ClearAndUpdate()
        {
            try
            {
                string updSql = GetUpdateSql();
                string clearSql = GetClearSql();
                if (updSql.Length == 0) return -1001;
                //
                using(SqlConnection cnn = new SqlConnection(CnnStr))
                {
                    cnn.Open();
                    SqlCommand cmd = cnn.CreateCommand();
                    using (SqlTransaction tran = cnn.BeginTransaction())
                    {
                        cmd.Transaction = tran;
                        // 清理数据
                        cmd.CommandText = clearSql;
                        cmd.ExecuteNonQuery();
                        // 更新数据
                        cmd.CommandText = updSql;
                        long counter = 0;
                        for(int row = 0; row < Data.Rows.Count; row++)
                        {
                            cmd.Parameters.Clear();
                            for (int col = 0; col < Data.Columns.Count; col++)
                            {
                                cmd.Parameters.AddWithValue("@data" + col,
                                    Data.Rows[row][col]);
                            }
                            counter += cmd.ExecuteNonQuery();
                        }
                        tran.Commit();
                        return counter;
                    }
                }
            }
            catch { return -1000; }
        }
        //
    }
}

SQL Server数据库批量导入时,与MySQL数据库比较需要注意以下一些区别。

SQL Server数据库中的对象名使用一对方括号定义,如[cdb_cs1]、[t1]。

约定SQL Server数据库语句中的参数名称使用@符号定义,而MySQL语句中的参数名使用问号(?)定义。

SQL Server语句中,exists关键字与条件的应用相似,需要在where关键字后使用,如代码中使用的"select 1 where exists(查询语句)",当查询语句包含查询结果时,语句会返回1,在InsertOrUpdate()方法中会根据查询结果是否为1确定记录是否存在。

tBatch组件的使用

在应用中使用tBatch接口组件时,可以将其扩展到tDbFactory接口组件中,也可以使用独立的组件使用,可以根据需要灵活设置和应用。

下面的代码(cfx/data/tDbFactory.cs)修改了tDbFactory接口的实现基类的代码,将tBatch组件添加到数据库工厂组件中。

C#
using System.Data;
namespace cfx.data
{
    public interface tDbFactory
    {
        string CnnStr { get; }
        tInsert GetInsert(string table);
        tUpdate GetUpdate(string table);
        tDelete GetDelete(string table);
        tQuery GetQuery(string source);
        //
        tDbJet GetJet();
        //
        tBatch GetBatch(string table, DataTable data);
    }
    //
    public abstract class tDbFactoryBase : tDbFactory
    {
        protected string myCnnStr;
        public tDbFactoryBase(string cnnstr)
        {
            myCnnStr = cnnstr;
        }
        //
        public string CnnStr { get { return myCnnStr; } }
        //
        public abstract tInsert GetInsert(string table);
        public abstract tUpdate GetUpdate(string table);
        public abstract tDelete GetDelete(string table);
        public abstract tQuery GetQuery(string source);
        //
        public abstract tDbJet GetJet();
        //
        public abstract tBatch GetBatch(string table, DataTable data);
    }
}

然后,在实现tDbFactory接口的类型中添加GetBatch()方法,如下面的代码(cfx/data/mysql/tMySqlFactory.cs)就是MySQL数据库的工厂实现类。

C#
using System.Data;
namespace cfx.data.mysql
{
    public class tMySqlFactory : tDbFactoryBase
    {
        public tMySqlFactory(string cnnstr) : base(cnnstr) { }
        //
        public static tMySqlFactory Create(string cnnstr)
        {
            return new tMySqlFactory(cnnstr);
        }
        //
        public override tInsert GetInsert(string table)
        {
            return tMySqlInsert.Create(myCnnStr, table);
        }
        //
        public override tUpdate GetUpdate(string table)
        {
            return tMySqlUpdate.Create(myCnnStr, table);
        }
        //
        public override tDelete GetDelete(string table)
        {
            return tMySqlDelete.Create(myCnnStr, table);
        }
        //
        public override tQuery GetQuery(string source)
        {
            return tMySqlQuery.Create(myCnnStr, source);
        }
        //
        public override tDbJet GetJet()
        {
            return tMySqlJet.Create(myCnnStr);
        }
        //
        public override tBatch GetBatch(string table, DataTable data)
        {
            return tMySqlBatch.Create(myCnnStr, table, data);
        }
        //
    }
}

下面的代码(cfx/data/sql/tSqlFactory.cs)是修改后的SQL Server数据库工厂类。

C#
using System.Data;

namespace cfx.data.sql
{
    public class tSqlFactory : tDbFactoryBase
    {
        public tSqlFactory(string cnnstr) : base(cnnstr) { }
        //
        public static tSqlFactory Create(string cnnstr)
        {
            return new tSqlFactory(cnnstr);
        }
        //
        public override tInsert GetInsert(string table)
        {
            return tSqlInsert.Create(myCnnStr, table);
        }
        //
        public override tUpdate GetUpdate(string table)
        {
            return tSqlUpdate.Create(myCnnStr, table);
        }
        //
        public override tDelete GetDelete(string table)
        {
            return tSqlDelete.Create(myCnnStr, table);
        }
        //
        public override tQuery GetQuery(string source)
        {
            return tSqlQuery.Create(myCnnStr, source);
        }
        //
        public override tDbJet GetJet()
        {
            return tSqlJet.Create(myCnnStr);
        }
        //
        public override tBatch GetBatch(string table, DataTable data)
        {
            return tSqlBatch.Create(myCnnStr, table, data);
        }
        //
    }
}

如果tApp中定义了dbFact1字段为数据库工厂对象,则在应用中可以参考如下代码使用tBatch组件。

C#
using System;
using System.Data;
using cfx.office;

namespace csfx_demo
{
    class Program
    {
        static void Main(string[] args)
        {
            DataTable data = tExcelReader.Create(@"d:\t1.xls").Read();
            long result = tApp.dbFact1.GetBatch("t1", data).Insert();
            Console.WriteLine(result);
        }
    }
}