You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

SSIS脚本需求:按SQL列值分组并按200条记录拆分文件

SSIS C#脚本改造:按记录数拆分输出文件

以下是改造后的完整脚本,实现当每组记录达到指定数量(如200条)时自动拆分文件,文件名格式为nameoffile_columnvalue_datetime_序号.txt:

public void Main()
{
    string datetime_1 = DateTime.Now.ToString("dd/MM/yyyy");
    string datetime = datetime_1.Replace("/", String.Empty);

    try
    {
        // 声明变量
        string FileNamePart = Dts.Variables["User::FlatFileNamePart"].Value.ToString();
        string DestinationFolder = Dts.Variables["User::DestinationFolder"].Value.ToString();
        string TableName = Dts.Variables["User::TableName"].Value.ToString();
        string FileDelimiter = Dts.Variables["User::FileDelimiter"].Value.ToString();
        string FileExtension = Dts.Variables["User::FileExtension"].Value.ToString();
        string ColumnNameForGrouping = Dts.Variables["User::ColumnNameForGrouping"].Value.ToString();
        string SubFolder = Dts.Variables["User::SubFolder"].Value.ToString();
        Int32 RecordCntPerFile = (Int32)Dts.Variables["User::RecordsPerFile"].Value;
        string RecordCntPerFileDecimal = RecordCntPerFile + ".0";

        // 从SSIS包获取ADO.NET连接
        SqlConnection myADONETConnection = (SqlConnection)(Dts.Connections["AR_GSLO_OLTP"].AcquireConnection(Dts.Transaction) as SqlConnection);

        // 获取分组列的所有唯一值
        string query = $"Select distinct {ColumnNameForGrouping} from {TableName}";
        SqlCommand cmd = new SqlCommand(query, myADONETConnection);
        DataTable dt = new DataTable();
        dt.Load(cmd.ExecuteReader());
        myADONETConnection.Close();

        // 遍历每个分组值
        foreach (DataRow dt_row in dt.Rows)
        {
            string ColumnValue = dt_row[0].ToString();

            // 加载当前分组的所有数据
            string queryString = $"SELECT * from {TableName} Where {ColumnNameForGrouping}='{ColumnValue}'";
            SqlDataAdapter adapter = new SqlDataAdapter(queryString, myADONETConnection);
            DataSet ds = new DataSet();
            adapter.Fill(ds);

            foreach (DataTable d_table in ds.Tables)
            {
                int ColumnCount = d_table.Columns.Count;
                int currentRecordCount = 0;
                int fileIndex = 1;
                StreamWriter sw = null;

                // 遍历每一行数据
                foreach (DataRow dr in d_table.Rows)
                {
                    // 当记录数达到阈值或为第一条记录时,创建新文件
                    if (currentRecordCount % RecordCntPerFile == 0)
                    {
                        sw?.Close(); // 关闭之前的文件流
                        string FileFullPath = $@"{SubFolder}\{FileNamePart}_{ColumnValue}_{datetime}_{fileIndex}{FileExtension}";
                        sw = new StreamWriter(FileFullPath, false);

                        // 若需要写入表头,取消下面的注释
                        //for (int ic = 0; ic < ColumnCount; ic++)
                        //{
                        //    sw.Write(d_table.Columns[ic]);
                        //    if (ic < ColumnCount - 1)
                        //    {
                        //        sw.Write(FileDelimiter);
                        //    }
                        //}
                        //sw.Write(sw.NewLine);

                        fileIndex++;
                    }

                    // 写入当前行数据
                    for (int ir = 0; ir < ColumnCount; ir++)
                    {
                        if (!Convert.IsDBNull(dr[ir]))
                        {
                            sw.Write(dr[ir].ToString());
                        }
                        if (ir < ColumnCount - 1)
                        {
                            sw.Write(FileDelimiter);
                        }
                    }
                    sw.Write(sw.NewLine);
                    currentRecordCount++;
                }

                sw?.Close(); // 关闭最后一个文件流
            }
        }

        Dts.TaskResult = (int)ScriptResults.Success;
    }
    catch (Exception exception)
    {
        // 创建错误日志文件
        using (StreamWriter sw = File.CreateText($@"{Dts.Variables["User::LogFolder"].Value.ToString()}\ErrorLog_{datetime}.log"))
        {
            sw.WriteLine(exception.ToString());
            Dts.TaskResult = (int)ScriptResults.Failure;
        }
    }
}

关键修改点

  • 新增currentRecordCount变量跟踪当前已写入的记录数,fileIndex变量标记当前文件序号
  • 在遍历数据行时,判断当前记录数是否达到阈值,达到则关闭当前文件流,创建新的文件并更新序号
  • 优化了字符串拼接方式,使用插值字符串提高可读性
  • 确保最后一个文件流被正确关闭
  • 补充了任务成功的返回标记Dts.TaskResult = (int)ScriptResults.Success

内容的提问来源于stack exchange,提问作者asdaswsad

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.19 05:05:31