-
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathCmdletInvokeDbaXQueryStream.cs
More file actions
136 lines (125 loc) · 6.72 KB
/
Copy pathCmdletInvokeDbaXQueryStream.cs
File metadata and controls
136 lines (125 loc) · 6.72 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
namespace DBAClientX.PowerShell;
/// <summary>Streams query rows through a DbaClientX provider.</summary>
/// <example>
/// <summary>Stream SQL Server rows.</summary>
/// <prefix>PS> </prefix>
/// <code>Invoke-DbaXQueryStream -Provider SqlServer -ConnectionString $connectionString -Query 'SELECT name FROM sys.databases' -ReturnType PSObject</code>
/// <para>Streams query rows without buffering the full result set.</para>
/// </example>
[Cmdlet(VerbsLifecycle.Invoke, "DbaXQueryStream", SupportsShouldProcess = true)]
[CmdletBinding()]
public sealed class CmdletInvokeDbaXQueryStream : AsyncPSCmdlet
{
/// <summary>Provider used to stream query rows.</summary>
[Parameter(Mandatory = true)]
public DbaXProvider Provider { get; set; }
/// <summary>Provider connection string, or SQLite database path.</summary>
[Parameter(Mandatory = true)]
[ValidateNotNullOrEmpty]
public string ConnectionString { get; set; } = string.Empty;
/// <summary>Query text to execute.</summary>
[Parameter(Mandatory = true)]
[ValidateNotNullOrEmpty]
public string Query { get; set; } = string.Empty;
/// <summary>Optional query parameters.</summary>
[Parameter(Mandatory = false)]
public Hashtable? Parameters { get; set; }
/// <summary>Executes through an active provider transaction when supported.</summary>
[Parameter(Mandatory = false)]
public SwitchParameter UseTransaction { get; set; }
/// <summary>Controls the returned object format. Defaults to PSObject so an ordinary PowerShell query emits every row.</summary>
[Parameter(Mandatory = false)]
[Alias("As")]
public ReturnType ReturnType { get; set; } = ReturnType.PSObject;
/// <summary>Optional command timeout in seconds.</summary>
[Parameter(Mandatory = false)]
public int QueryTimeout { get; set; }
/// <inheritdoc />
protected override async Task ProcessRecordAsync()
{
PowerShellHelpers.RejectFullConnectionTransactionSwitch(UseTransaction, MyInvocation.MyCommand.Name);
if (!ShouldProcess(Provider.ToString(), $"Stream {Provider} query"))
{
return;
}
#if NETSTANDARD2_1_OR_GREATER || NETCOREAPP3_0_OR_GREATER
var parameters = PowerShellHelpers.ToDictionaryOrNull(Parameters);
switch (Provider)
{
case DbaXProvider.SqlServer:
using (var client = new DBAClientX.SqlServer { ReturnType = ReturnType, CommandTimeout = QueryTimeout })
{
if (RequiresBufferedAggregate())
{
DbaXResultWriter.WriteResult(await client.QueryAsync(ConnectionString, Query, parameters, UseTransaction.IsPresent, CancelToken).ConfigureAwait(false), ReturnType, WriteObject);
}
else
{
await DbaXResultWriter.WriteRowsAsync(client.QueryStreamAsync(ConnectionString, Query, parameters, UseTransaction.IsPresent, CancelToken), ReturnType, WriteObject).ConfigureAwait(false);
}
}
break;
case DbaXProvider.PostgreSql:
using (var client = new DBAClientX.PostgreSql { ReturnType = ReturnType, CommandTimeout = QueryTimeout })
{
if (RequiresBufferedAggregate())
{
DbaXResultWriter.WriteResult(await client.QueryAsync(ConnectionString, Query, parameters, UseTransaction.IsPresent, CancelToken).ConfigureAwait(false), ReturnType, WriteObject);
}
else
{
await DbaXResultWriter.WriteRowsAsync(client.QueryStreamAsync(ConnectionString, Query, parameters, UseTransaction.IsPresent, CancelToken), ReturnType, WriteObject).ConfigureAwait(false);
}
}
break;
case DbaXProvider.MySql:
using (var client = new DBAClientX.MySql { ReturnType = ReturnType, CommandTimeout = QueryTimeout })
{
if (RequiresBufferedAggregate())
{
DbaXResultWriter.WriteResult(await client.QueryAsync(ConnectionString, Query, parameters, UseTransaction.IsPresent, CancelToken).ConfigureAwait(false), ReturnType, WriteObject);
}
else
{
await DbaXResultWriter.WriteRowsAsync(client.QueryStreamAsync(ConnectionString, Query, parameters, UseTransaction.IsPresent, CancelToken), ReturnType, WriteObject).ConfigureAwait(false);
}
}
break;
case DbaXProvider.Oracle:
using (var client = new DBAClientX.Oracle { ReturnType = ReturnType, CommandTimeout = QueryTimeout })
{
if (RequiresBufferedAggregate())
{
DbaXResultWriter.WriteResult(await client.QueryAsync(ConnectionString, Query, parameters, UseTransaction.IsPresent, CancelToken).ConfigureAwait(false), ReturnType, WriteObject);
}
else
{
await DbaXResultWriter.WriteRowsAsync(client.QueryStreamAsync(ConnectionString, Query, parameters, UseTransaction.IsPresent, CancelToken), ReturnType, WriteObject).ConfigureAwait(false);
}
}
break;
case DbaXProvider.SQLite:
using (var client = new DBAClientX.SQLite { ReturnType = ReturnType, CommandTimeout = QueryTimeout })
{
var connectionString = DbaXProviderHelpers.GetValidatedSQLiteConnectionString(ConnectionString);
if (RequiresBufferedAggregate())
{
DbaXResultWriter.WriteResult(await client.QueryWithConnectionStringAsync(connectionString, Query, parameters, UseTransaction.IsPresent, CancelToken).ConfigureAwait(false), ReturnType, WriteObject);
}
else
{
await DbaXResultWriter.WriteRowsAsync(client.QueryStreamWithConnectionStringAsync(connectionString, Query, parameters, UseTransaction.IsPresent, CancelToken), ReturnType, WriteObject).ConfigureAwait(false);
}
}
break;
default:
throw new PSArgumentException($"Provider '{Provider}' is not supported.", nameof(Provider));
}
#else
await Task.Yield();
throw new NotSupportedException("Streaming is not supported on this platform.");
#endif
}
private bool RequiresBufferedAggregate()
=> ReturnType == ReturnType.DataTable || ReturnType == ReturnType.DataSet;
}