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
相关产品推荐
相关产品推荐

