我的目标是从数据源检索数据,向其中添加一些元数据并将其插入到另一个目标。
目标的架构比源(计算列)多四列。
我正在使用SqlBulkCopy,它需要一个包含所有列(包括计算的 4 个)的阅读器。
有没有办法手动向 DataReader 添加列?或者如果不可能有什么替代数据插入?
我需要这样做,并且还能够基于其他列创建列,并使列的值取决于读取器的行索引。这个类是IWrappedReader从 Dapper 实现的,如果这对你有用的话(对我来说)。该类并不完整,因为我没有实现所有IDataRecord字段,但您可以看看我是如何做的IDataRecord.GetInt32以查看简单的模式。
/// <inheritdoc />
public class WrappedDataReader : IDataReader
{
private readonly IList<AdditionalField> _additionalFields;
private readonly int _originalOrdinalCount;
private IDbCommand _cmd;
private int _currentRowIndex = -1; //The first Read() will make this 0
private IDataReader _underlyingReader;
public WrappedDataReader(IDataReader underlyingReader, IList<AdditionalField> additionalFields)
{
_additionalFields = additionalFields;
_underlyingReader = underlyingReader;
var schema = Reader.GetSchemaTable();
if (schema == null)
{
throw new ObjectDisposedException(GetType().Name);
}
_originalOrdinalCount = schema.Rows.Count;
}
public object this[int i]
{
get { throw new NotImplementedException(); }
}
public IDataReader Reader
{
get
{
if (_underlyingReader == null)
{
throw new ObjectDisposedException(GetType().Name);
}
return _underlyingReader;
}
}
IDbCommand IWrappedDataReader.Command
{
get
{
if (_cmd == null)
{
throw new ObjectDisposedException(GetType().Name);
}
return _cmd;
}
}
void IDataReader.Close() => _underlyingReader?.Close();
int IDataReader.Depth => Reader.Depth;
DataTable IDataReader.GetSchemaTable()
{
var rv = Reader.GetSchemaTable();
if (rv == null)
{
throw new ObjectDisposedException(GetType().Name);
}
for (var i = 0; i < _additionalFields.Count; i++)
{
var row = rv.NewRow();
row["ColumnName"] = _additionalFields[i].ColumnName;
row["ColumnOrdinal"] = GetAppendColumnOrdinal(i);
row["DataType"] = _additionalFields[i].DataType;
rv.Rows.Add(row);
}
return rv;
}
bool IDataReader.IsClosed => _underlyingReader?.IsClosed ?? true;
bool IDataReader.NextResult() => Reader.NextResult();
bool IDataReader.Read()
{
_currentRowIndex++;
return Reader.Read();
}
int IDataReader.RecordsAffected => Reader.RecordsAffected;
void IDisposable.Dispose()
{
_underlyingReader?.Close();
_underlyingReader?.Dispose();
_underlyingReader = null;
_cmd?.Dispose();
_cmd = null;
}
int IDataRecord.FieldCount => Reader.FieldCount + _additionalFields.Count;
bool IDataRecord.GetBoolean(int i) => Reader.GetBoolean(i);
byte IDataRecord.GetByte(int i) => Reader.GetByte(i);
long IDataRecord.GetBytes(int i, long fieldOffset, byte[] buffer, int bufferoffset, int length) =>
Reader.GetBytes(i, fieldOffset, buffer, bufferoffset, length);
char IDataRecord.GetChar(int i) => Reader.GetChar(i);
long IDataRecord.GetChars(int i, long fieldoffset, char[] buffer, int bufferoffset, int length) =>
Reader.GetChars(i, fieldoffset, buffer, bufferoffset, length);
IDataReader IDataRecord.GetData(int i) => Reader.GetData(i);
string IDataRecord.GetDataTypeName(int i) => Reader.GetDataTypeName(i);
DateTime IDataRecord.GetDateTime(int i) => Reader.GetDateTime(i);
decimal IDataRecord.GetDecimal(int i) => Reader.GetDecimal(i);
double IDataRecord.GetDouble(int i) => Reader.GetDouble(i);
Type IDataRecord.GetFieldType(int i) => Reader.GetFieldType(i);
float IDataRecord.GetFloat(int i) => Reader.GetFloat(i);
Guid IDataRecord.GetGuid(int i) => Reader.GetGuid(i);
short IDataRecord.GetInt16(int i) => Reader.GetInt16(i);
int IDataRecord.GetInt32(int i)
{
return i >= _originalOrdinalCount ? (int) ExecuteAdditionalFieldFunc(i) : Reader.GetInt32(i);
}
long IDataRecord.GetInt64(int i) => Reader.GetInt64(i);
string IDataRecord.GetName(int i)
{
return i >= _originalOrdinalCount ? _additionalFields[GetAppendColumnIndex(i)].ColumnName : Reader.GetName(i);
}
int IDataRecord.GetOrdinal(string name)
{
for (var i = 0; i < _additionalFields.Count; i++)
{
if (name.Equals(_additionalFields[i].ColumnName, StringComparison.OrdinalIgnoreCase))
{
return GetAppendColumnOrdinal(i);
}
}
return Reader.GetOrdinal(name);
}
string IDataRecord.GetString(int i) => Reader.GetString(i);
object IDataRecord.GetValue(int i)
{
return i >= _originalOrdinalCount ? ExecuteAdditionalFieldFunc(i) : Reader.GetValue(i);
}
int IDataRecord.GetValues(object[] values) => Reader.GetValues(values);
bool IDataRecord.IsDBNull(int i)
{
return i >= _originalOrdinalCount ? ExecuteAdditionalFieldFunc(i) == null : Reader.IsDBNull(i);
}
object IDataRecord.this[string name]
{
get
{
var ordinal = ((IDataRecord) this).GetOrdinal(name);
return ((IDataRecord) this).GetValue(ordinal);
}
}
object IDataRecord.this[int i] => ((IDataRecord) this).GetValue(i);
private int GetAppendColumnOrdinal(int index)
{
return _originalOrdinalCount + index;
}
private int GetAppendColumnIndex(int oridinal)
{
return oridinal - _originalOrdinalCount;
}
private object ExecuteAdditionalFieldFunc(int oridinal)
{
return _additionalFields[GetAppendColumnIndex(oridinal)].Func(_currentRowIndex, Reader);
}
public struct AdditionalField
{
public AdditionalField(string columnName, Type dataType, Func<int, IDataReader, object> func = null)
{
ColumnName = columnName;
DataType = dataType;
Func = func;
}
public string ColumnName;
public Type DataType;
public Func<int, IDataReader, object> Func;
}
}
Run Code Online (Sandbox Code Playgroud)
这是我编写的一个快速 NUnit 测试,用于展示其用法
[Test]
public void ReaderTest()
{
using (var conn = new SqlConnection(ConnectionSettingsCollection.Default))
{
conn.Open();
const string sql = @"
SELECT 1 as OriginalField
UNION
SELECT -500 as OriginalField
UNION
SELECT 100 as OriginalField
";
var additionalFields = new[]
{
new WrappedDataReader.AdditionalField("StaticField", typeof(int), delegate { return "X"; }),
new WrappedDataReader.AdditionalField("CounterField", typeof(int), (i, reader) => i),
new WrappedDataReader.AdditionalField("ComputedField", typeof(int), (i, reader) => (int) reader["OriginalField"] + 1000)
};
const string expectedJson = @"
[
{""OriginalField"":-500,""StaticField"":""X"",""CounterField"":0,""ComputedField"":500},
{""OriginalField"":1, ""StaticField"":""X"",""CounterField"":1,""ComputedField"":1001},
{""OriginalField"":100, ""StaticField"":""X"",""CounterField"":2,""ComputedField"":1100}
]
";
var actualJson = ToJson(new WrappedDataReader(new SqlCommand(sql, conn).ExecuteReader(), additionalFields));
Assert.Zero(CultureInfo.InvariantCulture.CompareInfo.Compare(expectedJson, actualJson, CompareOptions.IgnoreSymbols));
}
}
private static string ToJson(IDataReader reader)
{
using (var strWriter = new StringWriter(new StringBuilder()))
using (var jsonWriter = new JsonTextWriter(strWriter))
{
jsonWriter.WriteStartArray();
while (reader.Read())
{
jsonWriter.WriteStartObject();
for (var i = 0; i < reader.FieldCount; i++)
{
jsonWriter.WritePropertyName(reader.GetName(i));
jsonWriter.WriteValue(reader[i]);
}
jsonWriter.WriteEndObject();
}
jsonWriter.WriteEndArray();
return strWriter.ToString();
}
}
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
6107 次 |
| 最近记录: |