From df60067b6e7803236840a895f2e316f620d70600 Mon Sep 17 00:00:00 2001 From: Shinai Yang Date: Wed, 16 Sep 2026 16:55:23 +0800 Subject: [PATCH 1/4] Add Spark 4.0 Driver API coverage and .NET Interactive support --- azure-pipelines-e2e-tests-template.yml | 34 +++- azure-pipelines-pr.yml | 2 +- .../InteractiveTests.cs | 122 ++++++++++++ ...tensions.DotNet.Interactive.E2ETest.csproj | 10 + .../AssemblyKernelExtensionTests.cs | 101 ++++++++++ .../AssemblyKernelExtension.cs | 30 ++- .../ReferencedPackagesExtractor.cs | 6 +- .../IpcTests/Sql/CatalogTests.cs | 1 + .../IpcTests/Sql/ColumnTests.cs | 1 + .../IpcTests/Sql/DataFrameFunctionsTests.cs | 1 + .../IpcTests/Sql/DataFrameReaderTests.cs | 1 + .../IpcTests/Sql/DataFrameTests.cs | 3 +- .../IpcTests/Sql/DataFrameWriterTests.cs | 1 + .../IpcTests/Sql/DataFrameWriterV2Tests.cs | 44 ++--- .../Sql/DriverApiCompatibilityTests.cs | 186 ++++++++++++++++++ .../Sql/Expressions/WindowSpecTests.cs | 1 + .../IpcTests/Sql/Expressions/WindowTests.cs | 1 + .../IpcTests/Sql/FunctionsTests.cs | 1 + .../IpcTests/Sql/RowTests.cs | 1 + .../IpcTests/Sql/RuntimeConfigTests.cs | 1 + .../Sql/SparkSessionExtensionsTests.cs | 1 + .../IpcTests/Sql/SparkSessionTests.cs | 1 + .../IpcTests/Sql/TypesTests.cs | 1 + src/csharp/Microsoft.Spark.sln | 7 + .../spark/sql/api/dotnet/SQLUtils.scala | 10 + .../spark/sql/api/dotnet/SQLUtilsTest.scala | 49 ++++- 26 files changed, 577 insertions(+), 40 deletions(-) create mode 100644 src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/InteractiveTests.cs create mode 100644 src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest.csproj create mode 100644 src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.UnitTest/AssemblyKernelExtensionTests.cs create mode 100644 src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DriverApiCompatibilityTests.cs diff --git a/azure-pipelines-e2e-tests-template.yml b/azure-pipelines-e2e-tests-template.yml index b91c614b9..9f7a7a4a8 100644 --- a/azure-pipelines-e2e-tests-template.yml +++ b/azure-pipelines-e2e-tests-template.yml @@ -268,12 +268,14 @@ stages: displayName: 'E2E tests' inputs: command: test - # Spark 4 covers the bridge, row retrieval, UDF, RDD, Broadcast and Worker lifecycle. Keep - # extension projects out of this lane until their Spark 4 compatibility is implemented. + # Spark 4 uses the qualified core categories. Interactive has its own + # testhost below because it initializes the bridge in REPL mode. ${{ if eq(split(test.version, '.')[0], '4') }}: projects: '**/Microsoft.Spark.E2ETest/Microsoft.Spark.E2ETest.csproj' ${{ else }}: - projects: '**/Microsoft.Spark*.E2ETest/*.csproj' + projects: | + **/Microsoft.Spark*.E2ETest/*.csproj + !**/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/*.csproj arguments: '--configuration $(buildConfiguration) ${{ option.testOptions }}' workingDirectory: $(Build.SourcesDirectory)$(PATH_SEPARATOR)dotnet-spark env: @@ -287,6 +289,24 @@ stages: DOTNET_SPARKFIXTURE_IVY_SETTINGS: $(SPARK_IVY_SETTINGS) DOTNET_SPARKFIXTURE_REPOSITORIES: ${{ parameters.managedOSSMavenFeedUrl }} + - ${{ if or(eq(split(test.version, '.')[0], '4'), eq(test.version, '3.5.3')) }}: + - task: DotNetCoreCLI@2 + displayName: '.NET Interactive cross-cell UDF E2E tests' + inputs: + command: test + projects: '**/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/*.csproj' + arguments: '--configuration $(buildConfiguration) --blame-hang --blame-hang-timeout 3min' + workingDirectory: $(Build.SourcesDirectory)$(PATH_SEPARATOR)dotnet-spark + env: + HADOOP_HOME: $(Build.BinariesDirectory)$(PATH_SEPARATOR)hadoop + SPARK_HOME: $(Build.BinariesDirectory)$(PATH_SEPARATOR)spark-${{ test.version }}-bin-hadoop + DOTNET_WORKER_DIR: $(CURRENT_DOTNET_WORKER_DIR) + DOTNET_SPARKFIXTURE_EXPECTED_SPARK_VERSION: ${{ test.version }} + DOTNET_ASSEMBLY_SEARCH_PATHS: $(Build.SourcesDirectory)$(PATH_SEPARATOR)dotnet-spark$(PATH_SEPARATOR)artifacts$(PATH_SEPARATOR)bin$(PATH_SEPARATOR)Microsoft.Spark.E2ETest$(PATH_SEPARATOR)$(buildConfiguration)$(PATH_SEPARATOR)net8.0 + ${{ if and(ne(parameters.managedOSSMavenFeed, ''), ne(parameters.managedOSSMavenFeedUrl, '')) }}: + DOTNET_SPARKFIXTURE_IVY_SETTINGS: $(SPARK_IVY_SETTINGS) + DOTNET_SPARKFIXTURE_REPOSITORIES: ${{ parameters.managedOSSMavenFeedUrl }} + # IO encryption is a SparkEnv startup setting. A separate test process is # required; modifying SparkContext.GetConf() in a test does not enable it. - ${{ if or(eq(split(test.version, '.')[0], '4'), eq(test.version, '3.5.3')) }}: @@ -350,7 +370,9 @@ stages: condition: ${{ test.enableBackwardCompatibleTests }} inputs: command: test - projects: '**/Microsoft.Spark*.E2ETest/*.csproj' + projects: | + **/Microsoft.Spark*.E2ETest/*.csproj + !**/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/*.csproj arguments: '--configuration $(buildConfiguration) ${{ option.backwardCompatibleTestOptions }}' workingDirectory: $(Build.SourcesDirectory)$(PATH_SEPARATOR)dotnet-spark env: @@ -377,7 +399,9 @@ stages: condition: ${{ test.enableForwardCompatibleTests }} inputs: command: test - projects: '**/Microsoft.Spark*.E2ETest/*.csproj' + projects: | + **/Microsoft.Spark*.E2ETest/*.csproj + !**/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/*.csproj arguments: '--configuration $(buildConfiguration) ${{ option.forwardCompatibleTestOptions }}' workingDirectory: $(Build.SourcesDirectory)$(PATH_SEPARATOR)dotnet-spark-${{ parameters.forwardCompatibleRelease }} env: diff --git a/azure-pipelines-pr.yml b/azure-pipelines-pr.yml index 4fac8ebdb..823b8d519 100644 --- a/azure-pipelines-pr.yml +++ b/azure-pipelines-pr.yml @@ -226,7 +226,7 @@ extends: - ${{ each pool in parameters.listOfE2ETestsPoolTypes }}: - pool: ${{ pool }} ${{ if startsWith(version, '4.0.') }}: - testOptions: '--filter "Category=Spark40Compatibility|Category=DataFrameRowRetrieval|Category=Broadcast|Category=ArrowUdf|Category=Avro|Category=Streaming"' + testOptions: '--filter "Category=Spark40Compatibility|Category=DataFrameRowRetrieval|Category=Broadcast|Category=ArrowUdf|Category=Avro|Category=Streaming|Category=DriverApi"' ${{ else }}: testOptions: "" backwardCompatibleTestOptions: $(backwardCompatibleTestOptions_${{ pool }}_${{ split(version, '.')[0] }}_${{ split(version, '.')[1] }}) diff --git a/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/InteractiveTests.cs b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/InteractiveTests.cs new file mode 100644 index 000000000..d4b990108 --- /dev/null +++ b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/InteractiveTests.cs @@ -0,0 +1,122 @@ +// Licensed to the .NET Foundation under one or more agreements. +// The .NET Foundation licenses this file to you under the MIT license. +// See the LICENSE file in the project root for more information. + +using System; +using System.Linq; +using System.Threading.Tasks; +using Microsoft.DotNet.Interactive; +using Microsoft.DotNet.Interactive.Commands; +using Microsoft.DotNet.Interactive.CSharp; +using Microsoft.DotNet.Interactive.Events; +using Microsoft.Spark.E2ETest; +using Microsoft.Spark.Sql; +using Xunit; + +namespace Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest +{ + public class InteractiveTests + { + [Fact] + public async Task TestCrossCellUdfAndRecovery() + { + string previousRepl = Environment.GetEnvironmentVariable("DOTNET_SPARK_RUNNING_REPL"); + string previousRunMode = Environment.GetEnvironmentVariable("SPARK_NET_RUN_MODE"); + try + { + // JvmBridge caches the REPL mode at construction. This project runs + // in a separate testhost, before any ordinary SQL fixture is created. + Environment.SetEnvironmentVariable("DOTNET_SPARK_RUNNING_REPL", "true"); + Environment.SetEnvironmentVariable("SPARK_NET_RUN_MODE", "N"); + using var fixture = new SparkFixture(); + using var kernel = new CompositeKernel(); + var csharp = new CSharpKernel(); + csharp.AddAssemblyReferences(new[] { typeof(SparkSession).Assembly.Location }); + kernel.Add(csharp); + await new AssemblyKernelExtension().OnLoadAsync(kernel); + + await SubmitAsync(kernel, @" + using System; + using System.Linq; + using Microsoft.Spark.Sql; + using static Microsoft.Spark.Sql.Functions; + public class CellOffset + { + public long Value; + public long Apply(long value) => value + Value; + } + var add = Udf(new CellOffset { Value = 7 }.Apply); + "); + await SubmitAsync(kernel, @" + var spark = SparkSession.Builder().GetOrCreate(); + var values = spark.Range(0, 3, 1, 2).Select(add(Col(""id""))) + .Collect().Select(row => row.GetAs(0)).OrderBy(value => value).ToArray(); + "); + AssertValues(csharp, "values", 7, 8, 9); + + await SubmitAsync(kernel, @" + public class LaterCell + { + public static long Apply(long value) => value * 3; + public static long Fail(long value) => throw new InvalidOperationException(""wi10-worker-failure""); + } + var multiply = Udf(LaterCell.Apply); + var fail = Udf(LaterCell.Fail); + "); + await SubmitAsync(kernel, @" + var later = spark.Range(0, 3, 1, 2).Select(multiply(Col(""id""))) + .Collect().Select(row => row.GetAs(0)).OrderBy(value => value).ToArray(); + "); + AssertValues(csharp, "later", 0, 3, 6); + + KernelCommandResult failedAction = await kernel.SendAsync(new SubmitCode(@" + public class FailedCell + { + public static long Apply(long value) => value + 100; + } + var fromFailedCell = Udf(FailedCell.Apply); + spark.Range(1).Select(fail(Col(""id""))).Collect().ToArray(); + ", "csharp")); + CommandFailed workerFailure = Assert.Single(failedAction.Events.OfType()); + Assert.NotNull(workerFailure.Exception); + Assert.Contains("wi10-worker-failure", workerFailure.Exception.ToString()); + + await SubmitAsync(kernel, @" + var recovered = spark.Range(0, 3, 1, 2).Select(fromFailedCell(Col(""id""))) + .Collect().Select(row => row.GetAs(0)).OrderBy(value => value).ToArray(); + "); + AssertValues(csharp, "recovered", 100, 101, 102); + + KernelCommandResult compileError = await kernel.SendAsync(new SubmitCode("int broken = ;", "csharp")); + CommandFailed compilerFailure = Assert.Single(compileError.Events.OfType()); + Assert.Contains("CS1525", compilerFailure.Message); + Assert.DoesNotContain("duplicate assembly", compilerFailure.Message); + + await SubmitAsync(kernel, @" + var afterCompileError = spark.Range(0, 3, 1, 2).Select(multiply(Col(""id""))) + .Collect().Select(row => row.GetAs(0)).OrderBy(value => value).ToArray(); + "); + AssertValues(csharp, "afterCompileError", 0, 3, 6); + } + finally + { + Environment.SetEnvironmentVariable("SPARK_NET_RUN_MODE", previousRunMode); + Environment.SetEnvironmentVariable("DOTNET_SPARK_RUNNING_REPL", previousRepl); + } + } + + private static async Task SubmitAsync(CompositeKernel kernel, string code) + { + KernelCommandResult result = await kernel.SendAsync(new SubmitCode(code, "csharp")); + string[] failures = result.Events.OfType().Select(failure => failure.Message).ToArray(); + Assert.True(failures.Length == 0, string.Join(Environment.NewLine, failures)); + Assert.Single(result.Events.OfType()); + } + + private static void AssertValues(CSharpKernel kernel, string name, params long[] expected) + { + Assert.True(kernel.TryGetValue(name, out long[] values)); + Assert.Equal(expected, values); + } + } +} diff --git a/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest.csproj b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest.csproj new file mode 100644 index 000000000..b85e2bff5 --- /dev/null +++ b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest/Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest.csproj @@ -0,0 +1,10 @@ + + + net8.0 + + + + + + + diff --git a/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.UnitTest/AssemblyKernelExtensionTests.cs b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.UnitTest/AssemblyKernelExtensionTests.cs new file mode 100644 index 000000000..c587f25f5 --- /dev/null +++ b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive.UnitTest/AssemblyKernelExtensionTests.cs @@ -0,0 +1,101 @@ +// Licensed to the .NET Foundation under one or more agreements. +// The .NET Foundation licenses this file to you under the MIT license. +// See the LICENSE file in the project root for more information. + +using System; +using System.Linq; +using System.Threading.Tasks; +using Microsoft.DotNet.Interactive; +using Microsoft.DotNet.Interactive.Commands; +using Microsoft.DotNet.Interactive.CSharp; +using Microsoft.DotNet.Interactive.Events; +using Xunit; + +namespace Microsoft.Spark.Extensions.DotNet.Interactive.UnitTest +{ + public class AssemblyKernelExtensionTests + { + [Theory] + [InlineData("3.0.0")] + [InlineData("3.5.3")] + [InlineData("4.0.0")] + [InlineData("4.0.1")] + [InlineData("4.0.2")] + [InlineData("4.0.3")] + [InlineData("4.0.4")] + public void TestPathsWithSpaces(string version) + { + AssemblyKernelExtension.ValidatePath("packages with spaces/package.nupkg", new Version(version)); + } + + [Theory] + [InlineData("2.4.8")] + [InlineData("4.1.0")] + [InlineData("5.0.0")] + public void TestUnsupportedPathsWithSpaces(string version) + { + Assert.Throws(() => + AssemblyKernelExtension.ValidatePath("packages with spaces/package.nupkg", new Version(version))); + + AssemblyKernelExtension.ValidatePath("packages/package.nupkg", new Version(version)); + } + + [Fact] + public void TestKernelWithoutNuGetProvider() + { + using var kernel = new CSharpKernel(); + Assert.Empty(new ReferencedPackagesExtractor(kernel).ResolvedPackageReferences); + } + + [Fact] + public void TestMissingKernelIsNotTreatedAsAnEmptyPackageList() + { + Assert.ThrowsAny(() => + new ReferencedPackagesExtractor(null).ResolvedPackageReferences.ToArray()); + } + + [Fact] + public async Task TestKernelWithoutCSharpStillHandlesCommands() + { + string previousMode = Environment.GetEnvironmentVariable("DOTNET_SPARK_RUNNING_REPL"); + try + { + using var kernel = new CompositeKernel { new NonCSharpKernel() }; + await new AssemblyKernelExtension().OnLoadAsync(kernel); + KernelCommandResult result = await kernel.SendAsync(new SubmitCode("hello", "other")); + Assert.Empty(result.Events.OfType()); + Assert.Single(result.Events.OfType()); + } + finally + { + Environment.SetEnvironmentVariable("DOTNET_SPARK_RUNNING_REPL", previousMode); + } + } + + [Fact] + public async Task TestInitialCompilationFailureDoesNotPublishAnAssembly() + { + string previousMode = Environment.GetEnvironmentVariable("DOTNET_SPARK_RUNNING_REPL"); + try + { + using var kernel = new CompositeKernel { new CSharpKernel() }; + await new AssemblyKernelExtension().OnLoadAsync(kernel); + KernelCommandResult result = await kernel.SendAsync(new SubmitCode("int broken = ;", "csharp")); + CommandFailed failure = Assert.Single(result.Events.OfType()); + Assert.Contains("CS1525", failure.Message); + } + finally + { + Environment.SetEnvironmentVariable("DOTNET_SPARK_RUNNING_REPL", previousMode); + } + } + + private sealed class NonCSharpKernel : Kernel + { + public NonCSharpKernel() : base("other") + { + RegisterCommandHandler((command, context) => Task.CompletedTask); + } + } + } +} diff --git a/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive/AssemblyKernelExtension.cs b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive/AssemblyKernelExtension.cs index 33b1e23a1..74e887fe7 100644 --- a/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive/AssemblyKernelExtension.cs +++ b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive/AssemblyKernelExtension.cs @@ -59,18 +59,21 @@ public Task OnLoadAsync(Kernel kernel) compositeKernel.AddMiddleware(async (command, context, next) => { + var previousScript = cSharpKernel?.ScriptState?.Script; await next(command, context); if ((context.HandlingKernel is CSharpKernel kernel) && (command is SubmitCode) && + kernel.ScriptState?.Script != null && + !ReferenceEquals(previousScript, kernel.ScriptState.Script) && TryGetSparkSession(out SparkSession sparkSession) && TryEmitAssembly(kernel, tempDir.FullName, out string assemblyPath)) { - sparkSession.SparkContext.AddFile(assemblyPath); + AddFile(sparkSession, assemblyPath); foreach (string filePath in GetPackageFiles(tempDir.FullName)) { - sparkSession.SparkContext.AddFile(filePath); + AddFile(sparkSession, filePath); } } }); @@ -79,6 +82,22 @@ public Task OnLoadAsync(Kernel kernel) return Task.CompletedTask; } + private static void AddFile(SparkSession sparkSession, string path) + { + Version version = SparkEnvironment.SparkVersion; + if (version.Major == 4 && version.Minor == 0) + { + // Spark 4 SQL tasks only receive files registered in their session's + // artifact scope. Establish that scope and add the file in one JVM call. + sparkSession.Reference.Jvm.CallStaticJavaMethod( + "org.apache.spark.sql.api.dotnet.SQLUtils", "addFile", sparkSession, path); + } + else + { + sparkSession.SparkContext.AddFile(path); + } + } + private DirectoryInfo CreateTempDirectory() { string envTempDir = Environment.GetEnvironmentVariable(TempDirEnvVar); @@ -147,15 +166,16 @@ private IEnumerable GetPackageFiles(string path) /// - https://github.com/apache/spark/pull/26773 /// /// The path to validate. - private void ValidatePath(string path) + /// The Spark version, or the active session version. + internal static void ValidatePath(string path, Version version = null) { if (!path.Contains(" ")) { return; } - Version version = SparkEnvironment.SparkVersion; - if (version.Major != 3) + version ??= SparkEnvironment.SparkVersion; + if (version.Major != 3 && (version.Major != 4 || version.Minor != 0)) { throw new NotSupportedException($"Spark {version} not supported."); } diff --git a/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive/ReferencedPackagesExtractor.cs b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive/ReferencedPackagesExtractor.cs index 6de4a21de..81928b6c4 100644 --- a/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive/ReferencedPackagesExtractor.cs +++ b/src/csharp/Extensions/Microsoft.Spark.Extensions.DotNet.Interactive/ReferencedPackagesExtractor.cs @@ -46,7 +46,8 @@ internal virtual IEnumerable ResolvedPackageReferences if (restoreContext is null) { - throw new Exception("PackageRestoreContext was not found in the kernel, try using older version of Dotnet.Interactive"); + // A CSharpKernel can run without the optional NuGet provider. + return Array.Empty(); } // Extract the ResolvedPackageReferences property (internal) @@ -54,7 +55,8 @@ internal virtual IEnumerable ResolvedPackageReferences .GetType() .GetProperty("ResolvedPackageReferences", BindingFlags.Instance | BindingFlags.NonPublic | BindingFlags.Public); - return resolvedPackagesProp?.GetValue(restoreContext) as IEnumerable; + return resolvedPackagesProp?.GetValue(restoreContext) as IEnumerable + ?? throw new Exception("Failed to retrieve referenced packages from kernel, try using older version of Dotnet.Interactive"); } } } diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/CatalogTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/CatalogTests.cs index 630fb3c54..5f73d2564 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/CatalogTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/CatalogTests.cs @@ -14,6 +14,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class CatalogTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/ColumnTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/ColumnTests.cs index 6ffa2b3c7..dce9de694 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/ColumnTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/ColumnTests.cs @@ -11,6 +11,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class ColumnTests { /// diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameFunctionsTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameFunctionsTests.cs index 8f004b4f8..a44d202fb 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameFunctionsTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameFunctionsTests.cs @@ -11,6 +11,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class DataFrameFunctionsTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameReaderTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameReaderTests.cs index feb9b33ff..802d414c7 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameReaderTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameReaderTests.cs @@ -10,6 +10,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class DataFrameReaderTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameTests.cs index b8a65e3d6..51a4cc247 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameTests.cs @@ -23,7 +23,8 @@ namespace Microsoft.Spark.E2ETest.IpcTests { - [Collection("Spark E2E Tests")] + [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class DataFrameTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterTests.cs index bd3ce5804..5ffd17270 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterTests.cs @@ -10,6 +10,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class DataFrameWriterTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterV2Tests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterV2Tests.cs index 88491017f..7f8ac7aa5 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterV2Tests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterV2Tests.cs @@ -11,6 +11,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class DataFrameWriterV2Tests { private readonly SparkSession _spark; @@ -50,30 +51,27 @@ public void TestSignaturesV3_0_X() Assert.IsType(dfwV2.PartitionedBy(df.Col("age"))); - // Throws the following exception: - // org.apache.spark.sql.AnalysisException: REPLACE TABLE AS SELECT is only supported - // with v2 tables. - Assert.Throws(() => dfwV2.Replace()); - - // Throws the following exception: - // org.apache.spark.sql.AnalysisException: REPLACE TABLE AS SELECT is only supported - // with v2 tables. - Assert.Throws(() => dfwV2.CreateOrReplace()); - - // Throws the following exception: - // org.apache.spark.sql.AnalysisException: Table default.testtable does not support - // append in batch mode. - Assert.Throws(() => dfwV2.Append()); - - // Throws the following exception: - // org.apache.spark.sql.AnalysisException: Table default.testtable does not support - // overwrite by filter in batch mode. - Assert.Throws(() => dfwV2.Overwrite(df.Col("age"))); + // The JSON provider creates a V1 table, which must reject these V2 operations. + AssertUnsupportedTableOperation(() => dfwV2.Replace()); + AssertUnsupportedTableOperation(() => dfwV2.CreateOrReplace()); + AssertUnsupportedTableOperation(() => dfwV2.Append()); + AssertUnsupportedTableOperation(() => dfwV2.Overwrite(df.Col("age"))); + AssertUnsupportedTableOperation(() => dfwV2.OverwritePartitions()); + } - // Throws the following exception: - // org.apache.spark.sql.AnalysisException: Table default.testtable does not support - // dynamic overwrite in batch mode. - Assert.Throws(() => dfwV2.OverwritePartitions()); + private static void AssertUnsupportedTableOperation(Action operation) + { + Exception exception = Assert.Throws(operation); + JvmException jvmException = Assert.IsType(exception.InnerException); + Assert.Contains("org.apache.spark.sql.AnalysisException", jvmException.Message); + + string message = jvmException.Message.ToLowerInvariant(); + Assert.True( + message.Contains("only supported") || + message.Contains("does not support") || + message.Contains("unsupported") || + message.Contains("cannot write into v1 table"), + jvmException.Message); } } } diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DriverApiCompatibilityTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DriverApiCompatibilityTests.cs new file mode 100644 index 000000000..b034f4427 --- /dev/null +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DriverApiCompatibilityTests.cs @@ -0,0 +1,186 @@ +// Licensed to the .NET Foundation under one or more agreements. +// The .NET Foundation licenses this file to you under the MIT license. +// See the LICENSE file in the project root for more information. + +using System; +using System.IO; +using System.Linq; +using Microsoft.Spark.Sql; +using Microsoft.Spark.Sql.Catalog; +using Microsoft.Spark.Sql.Types; +using Microsoft.Spark.UnitTest.TestUtils; +using Xunit; +using static Microsoft.Spark.Sql.Functions; + +namespace Microsoft.Spark.E2ETest.IpcTests +{ + [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] + public class DriverApiCompatibilityTests + { + private readonly SparkSession _spark; + + public DriverApiCompatibilityTests(SparkFixture fixture) + { + _spark = fixture.Spark; + } + + [Fact] + public void TestSqlDataFrameParquetRoundTrip() + { + DataFrame orders = _spark.Sql(@" + SELECT * FROM VALUES + (1, 'w', 2, 10), + (2, 'w', 3, 7), + (3, 'e', 4, 5), + (4, 'e', 1, 8), + (5, 'x', 5, 6) + AS orders(order_id, region_id, quantity, unit_price)"); + DataFrame regions = _spark.Sql(@" + SELECT * FROM VALUES ('w', 'West'), ('e', 'East') + AS regions(region_id, region)"); + + DataFrame totals = orders + .Select( + Col("region_id"), + (Col("quantity") * Col("unit_price")).Alias("amount")) + .Filter(Col("amount") >= 20) + .Join(regions, "region_id") + .GroupBy("region") + .Agg(Sum("amount").Alias("total"), Count(Lit(1)).Alias("order_count")); + + AssertOrderTotals(totals); + Assert.Equal(new[] { "region", "total", "order_count" }, totals.Columns()); + Assert.IsType(totals.Schema().Fields[0].DataType); + Assert.IsType(totals.Schema().Fields[1].DataType); + Assert.IsType(totals.Schema().Fields[2].DataType); + + using (var directory = new TemporaryDirectory()) + { + string path = Path.Combine(directory.Path, "orders"); + totals.Write().Parquet(path); + DataFrame restored = _spark.Read().Parquet(path); + + // Parquet reads make fields nullable even when their input was non-nullable. + var expectedSchema = new StructType(new[] + { + new StructField("region", new StringType()), + new StructField("total", new LongType()), + new StructField("order_count", new LongType()) + }); + Assert.Equal(expectedSchema, restored.Schema()); + AssertOrderTotals(restored); + } + } + + [Fact] + public void TestCatalogViewAndColumnResults() + { + string viewName = "driver_api_" + Guid.NewGuid().ToString("N"); + Catalog catalog = _spark.Catalog; + DataFrame result = _spark.Sql( + "SELECT 3 AS quantity, 7 AS unit_price, 'spark' AS product") + .Select( + (Col("quantity") * Col("unit_price")).Alias("amount"), + Upper(Col("product")).Alias("name"), + (Col("quantity") >= 2).Alias("bulk"), + Col("quantity").Cast("string").Alias("quantity_text")); + + try + { + result.CreateOrReplaceTempView(viewName); + Assert.True(catalog.TableExists(viewName)); + + Table table = catalog.GetTable(viewName); + Assert.Equal(viewName, table.Name); + Assert.True(table.IsTemporary); + Assert.Equal("TEMPORARY", table.TableType); + Assert.Null(table.Database); + Assert.Null(table.Description); + + Row listedTable = Assert.Single(catalog.ListTables() + .Filter(Col("name") == viewName) + .Select("name", "tableType", "isTemporary") + .Collect()); + Assert.Equal(new object[] { viewName, "TEMPORARY", true }, listedTable.Values); + + var expectedColumns = new[] + { + new object[] { "amount", "int", false, false }, + new object[] { "bulk", "boolean", false, false }, + new object[] { "name", "string", false, false }, + new object[] { "quantity_text", "string", false, false } + }; + Assert.Equal( + expectedColumns, + catalog.ListColumns(viewName) + .Select("name", "dataType", "isPartition", "isBucket") + .OrderBy("name") + .Collect() + .Select(row => row.Values)); + + Row row = Assert.Single(_spark.Table(viewName).Collect()); + Assert.Equal(new object[] { 21, "SPARK", true, "3" }, row.Values); + } + finally + { + catalog.DropTempView(viewName); + } + + Assert.False(catalog.TableExists(viewName)); + } + + [Fact] + public void TestAnsiCastFailureAndRecovery() + { + const string ansiKey = "spark.sql.ansi.enabled"; + const string invalidCast = "SELECT CAST('not-an-integer' AS INT) AS value"; + SparkSession session = _spark.NewSession(); + RuntimeConfig conf = session.Conf(); + string originalAnsi = conf.Get(ansiKey); + string sharedAnsi = _spark.Conf().Get(ansiKey); + + try + { + conf.Set(ansiKey, true); + Exception exception = Assert.Throws(() => + session.Sql(invalidCast).Collect().ToArray()); + JvmException jvmException = Assert.IsType(exception.InnerException); + Assert.Contains("NumberFormatException", jvmException.Message); + Assert.Contains("not-an-integer", jvmException.Message); + Assert.Equal(42, session.Sql("SELECT 40 + 2 AS value").First().GetAs(0)); + + conf.Set(ansiKey, false); + Row permissiveRow = Assert.Single(session.Sql(invalidCast).Collect()); + Assert.Null(permissiveRow.Values[0]); + } + finally + { + conf.Set(ansiKey, originalAnsi); + } + + Assert.Equal(originalAnsi, conf.Get(ansiKey)); + Assert.Equal(sharedAnsi, _spark.Conf().Get(ansiKey)); + Assert.Equal(42, session.Sql("SELECT 40 + 2 AS value").First().GetAs(0)); + } + + private static void AssertOrderTotals(DataFrame totals) + { + // GetAs handles small BIGINT values that arrive boxed as Int32 through Pickle. + Assert.Collection( + totals.OrderBy("region").Collect(), + row => + { + Assert.Equal("East", row.GetAs("region")); + Assert.Equal(20L, row.GetAs("total")); + Assert.Equal(1L, row.GetAs("order_count")); + }, + row => + { + Assert.Equal("West", row.GetAs("region")); + Assert.Equal(41L, row.GetAs("total")); + Assert.Equal(2L, row.GetAs("order_count")); + }); + } + } +} diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/Expressions/WindowSpecTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/Expressions/WindowSpecTests.cs index 62a6f6378..1cf873148 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/Expressions/WindowSpecTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/Expressions/WindowSpecTests.cs @@ -11,6 +11,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class WindowSpecTests { /// diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/Expressions/WindowTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/Expressions/WindowTests.cs index 885d1ff3c..8edd206b9 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/Expressions/WindowTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/Expressions/WindowTests.cs @@ -11,6 +11,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class WindowTests { /// diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/FunctionsTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/FunctionsTests.cs index fd3e02b9b..6f82d6299 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/FunctionsTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/FunctionsTests.cs @@ -13,6 +13,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class FunctionsTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/RowTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/RowTests.cs index 64c6d7155..0e4d69430 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/RowTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/RowTests.cs @@ -12,6 +12,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class RowTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/RuntimeConfigTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/RuntimeConfigTests.cs index 4994c2803..4048db637 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/RuntimeConfigTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/RuntimeConfigTests.cs @@ -8,6 +8,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class RuntimeConfigTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/SparkSessionExtensionsTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/SparkSessionExtensionsTests.cs index fef51c73c..8508aaa82 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/SparkSessionExtensionsTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/SparkSessionExtensionsTests.cs @@ -10,6 +10,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class SparkSessionExtensionsTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/SparkSessionTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/SparkSessionTests.cs index 3e2745b2c..ae42154c1 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/SparkSessionTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/SparkSessionTests.cs @@ -13,6 +13,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class SparkSessionTests { private readonly SparkSession _spark; diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/TypesTests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/TypesTests.cs index 12340febf..ba199eb61 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/TypesTests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/TypesTests.cs @@ -9,6 +9,7 @@ namespace Microsoft.Spark.E2ETest.IpcTests { [Collection("Spark E2E Tests")] + [Trait("Category", "DriverApi")] public class TypesTests { private readonly IJvmBridge _jvm; diff --git a/src/csharp/Microsoft.Spark.sln b/src/csharp/Microsoft.Spark.sln index b53c4de8e..050042112 100644 --- a/src/csharp/Microsoft.Spark.sln +++ b/src/csharp/Microsoft.Spark.sln @@ -35,6 +35,8 @@ Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Microsoft.Spark.Extensions. EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Microsoft.Spark.Extensions.DotNet.Interactive.UnitTest", "Extensions\Microsoft.Spark.Extensions.DotNet.Interactive.UnitTest\Microsoft.Spark.Extensions.DotNet.Interactive.UnitTest.csproj", "{7BDE09ED-04B3-41B2-A466-3D6F7225291E}" EndProject +Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest", "Extensions\Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest\Microsoft.Spark.Extensions.DotNet.Interactive.E2ETest.csproj", "{81252A97-8E87-4333-9BA7-0B9360FD21E6}" +EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Microsoft.Spark.Extensions.Hyperspace", "Extensions\Microsoft.Spark.Extensions.Hyperspace\Microsoft.Spark.Extensions.Hyperspace.csproj", "{70DDA4E9-1195-4A29-9AA1-96A8223A6D4F}" EndProject Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "Microsoft.Spark.Extensions.Hyperspace.E2ETest", "Extensions\Microsoft.Spark.Extensions.Hyperspace.E2ETest\Microsoft.Spark.Extensions.Hyperspace.E2ETest.csproj", "{C6019E44-C777-4DE2-B70E-EA025B7D044D}" @@ -93,6 +95,10 @@ Global {7BDE09ED-04B3-41B2-A466-3D6F7225291E}.Debug|Any CPU.Build.0 = Debug|Any CPU {7BDE09ED-04B3-41B2-A466-3D6F7225291E}.Release|Any CPU.ActiveCfg = Release|Any CPU {7BDE09ED-04B3-41B2-A466-3D6F7225291E}.Release|Any CPU.Build.0 = Release|Any CPU + {81252A97-8E87-4333-9BA7-0B9360FD21E6}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {81252A97-8E87-4333-9BA7-0B9360FD21E6}.Debug|Any CPU.Build.0 = Debug|Any CPU + {81252A97-8E87-4333-9BA7-0B9360FD21E6}.Release|Any CPU.ActiveCfg = Release|Any CPU + {81252A97-8E87-4333-9BA7-0B9360FD21E6}.Release|Any CPU.Build.0 = Release|Any CPU {70DDA4E9-1195-4A29-9AA1-96A8223A6D4F}.Debug|Any CPU.ActiveCfg = Debug|Any CPU {70DDA4E9-1195-4A29-9AA1-96A8223A6D4F}.Debug|Any CPU.Build.0 = Debug|Any CPU {70DDA4E9-1195-4A29-9AA1-96A8223A6D4F}.Release|Any CPU.ActiveCfg = Release|Any CPU @@ -112,6 +118,7 @@ Global {206E16CA-ED59-4F5E-8EA1-9BB7BEEACB63} = {71A19F75-8279-40AB-BEA0-7D4B153FC416} {9C32014D-8C0C-40F1-9ABA-C3BF19687E5C} = {71A19F75-8279-40AB-BEA0-7D4B153FC416} {7BDE09ED-04B3-41B2-A466-3D6F7225291E} = {71A19F75-8279-40AB-BEA0-7D4B153FC416} + {81252A97-8E87-4333-9BA7-0B9360FD21E6} = {71A19F75-8279-40AB-BEA0-7D4B153FC416} {70DDA4E9-1195-4A29-9AA1-96A8223A6D4F} = {71A19F75-8279-40AB-BEA0-7D4B153FC416} {C6019E44-C777-4DE2-B70E-EA025B7D044D} = {71A19F75-8279-40AB-BEA0-7D4B153FC416} EndGlobalSection diff --git a/src/scala/microsoft-spark-4-0/src/main/scala/org/apache/spark/sql/api/dotnet/SQLUtils.scala b/src/scala/microsoft-spark-4-0/src/main/scala/org/apache/spark/sql/api/dotnet/SQLUtils.scala index 08efc40cd..d2ae3c7aa 100644 --- a/src/scala/microsoft-spark-4-0/src/main/scala/org/apache/spark/sql/api/dotnet/SQLUtils.scala +++ b/src/scala/microsoft-spark-4-0/src/main/scala/org/apache/spark/sql/api/dotnet/SQLUtils.scala @@ -15,12 +15,22 @@ import org.apache.spark.SparkContext import org.apache.spark.api.python.{PythonAccumulatorV2, PythonBroadcast, PythonFunction, SimplePythonFunction} import org.apache.spark.broadcast.Broadcast import org.apache.spark.deploy.dotnet.DotnetRunner +import org.apache.spark.sql.classic.SparkSession import org.apache.spark.sql.execution.python.PythonUDFRunner import scala.util.control.NonFatal object SQLUtils { + /** + * Registers Interactive files in the same artifact scope used by the session's SQL queries. + */ + def addFile(session: SparkSession, path: String): Unit = { + session.artifactManager.withResources { + session.sparkContext.addFile(path) + } + } + /** * Exposes createPythonFunction to the .NET client to enable registering UDFs. */ diff --git a/src/scala/microsoft-spark-4-0/src/test/scala/org/apache/spark/sql/api/dotnet/SQLUtilsTest.scala b/src/scala/microsoft-spark-4-0/src/test/scala/org/apache/spark/sql/api/dotnet/SQLUtilsTest.scala index e66ff09e8..e8a378c2d 100644 --- a/src/scala/microsoft-spark-4-0/src/test/scala/org/apache/spark/sql/api/dotnet/SQLUtilsTest.scala +++ b/src/scala/microsoft-spark-4-0/src/test/scala/org/apache/spark/sql/api/dotnet/SQLUtilsTest.scala @@ -6,13 +6,15 @@ package org.apache.spark.sql.api.dotnet +import java.io.FileNotFoundException import java.net.URI -import java.nio.file.{Files, Path} +import java.nio.file.{Files, Path, Paths} -import org.apache.spark.SparkContext +import org.apache.spark.{JobArtifactSet, JobArtifactState, SparkConf, SparkContext, SparkFiles} import org.apache.spark.deploy.dotnet.DotnetRunner +import org.apache.spark.sql.classic.SparkSession import org.apache.spark.sql.execution.python.PythonUDFRunner -import org.junit.Assert.{assertEquals, assertThrows, assertTrue} +import org.junit.Assert.{assertEquals, assertFalse, assertThrows, assertTrue} import org.junit.Test @Test @@ -21,6 +23,47 @@ class SQLUtilsTest { private val EmptyFileSha256 = "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855" + @Test + def shouldAddFileToSessionAndRestoreArtifactState(): Unit = { + withTemporaryDirectory { directory => + val conf = new SparkConf(false) + .setMaster("local[1]") + .setAppName("SQLUtilsTest") + .set("spark.ui.enabled", "false") + .set("spark.sql.artifact.isolation.enabled", "true") + val sparkContext = new SparkContext(conf) + try { + val session = new SparkSession(sparkContext) + val file = Files.createFile(directory.resolve("submission.dll")) + val uri = file.toUri.toString + val originalState = JobArtifactSet.getCurrentJobArtifactState + val priorState = JobArtifactState("sql-utils-previous-session", None) + + JobArtifactSet.withActiveJobArtifactState(priorState) { + SQLUtils.addFile(session, uri) + + assertEquals(Some(priorState), JobArtifactSet.getCurrentJobArtifactState) + assertEquals(Set(uri), sparkContext.addedFiles(session.sessionUUID).keySet.toSet) + assertFalse(sparkContext.addedFiles.contains(priorState.uuid)) + assertFalse(sparkContext.addedFiles.get("default").exists(_.contains(uri))) + assertTrue(Files.isRegularFile(Paths.get( + SparkFiles.getRootDirectory(), session.sessionUUID, file.getFileName.toString))) + + assertThrows( + classOf[FileNotFoundException], + () => SQLUtils.addFile(session, directory.resolve("missing.dll").toUri.toString)) + + assertEquals(Some(priorState), JobArtifactSet.getCurrentJobArtifactState) + assertEquals(Set(uri), sparkContext.addedFiles(session.sessionUUID).keySet.toSet) + } + + assertEquals(originalState, JobArtifactSet.getCurrentJobArtifactState) + } finally { + sparkContext.stop() + } + } + } + @Test def shouldReturnRuntimeArtifactIdentityInFixedOrder(): Unit = { withTemporaryDirectory { directory => From a1d7cc806e1c89f70bdc6a9629763d1fbef478b4 Mon Sep 17 00:00:00 2001 From: Shinai Yang Date: Thu, 17 Sep 2026 11:54:38 +0800 Subject: [PATCH 2/4] Fix Windows Hadoop setup for Spark 4 bridge tests --- azure-pipelines-pr.yml | 31 +++++++++++++++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/azure-pipelines-pr.yml b/azure-pipelines-pr.yml index 823b8d519..b11342e9a 100644 --- a/azure-pipelines-pr.yml +++ b/azure-pipelines-pr.yml @@ -165,6 +165,34 @@ extends: mavenPomFile: src/scala/pom.xml options: '-Dmaven.resolver.transport=wagon' + - pwsh: | + $ErrorActionPreference = 'Stop' + # SparkContext.addFile needs winutils even in local-mode Scala tests on Windows. + $hadoopRoot = Join-Path "$(Build.BinariesDirectory)" 'hadoop' + $hadoopBin = Join-Path $hadoopRoot 'bin' + $archive = Join-Path "$(Build.BinariesDirectory)" 'hadoop-3.3.5.zip' + Invoke-WebRequest -Uri 'https://github.com/SparkSnail/winutils/releases/download/hadoop-3.3.5/hadoop-3.3.5.zip' -OutFile $archive + Expand-Archive -LiteralPath $archive -DestinationPath $hadoopRoot -Force + New-Item -ItemType Directory -Path $hadoopBin -Force | Out-Null + Copy-Item -LiteralPath (Join-Path $hadoopRoot 'hadoop-3.3.5\winutils.exe'), (Join-Path $hadoopRoot 'hadoop-3.3.5\hadoop.dll') -Destination $hadoopBin -Force + + # These Hadoop binaries depend on the VC++ 2010 x64 runtime, as in Windows E2E. + $runtimeDll = Join-Path $env:SystemRoot 'System32\MSVCR100.dll' + if (-not (Test-Path -LiteralPath $runtimeDll)) { + $installer = Join-Path "$(Build.BinariesDirectory)" 'vcredist_x64.exe' + Invoke-WebRequest -Uri 'https://download.microsoft.com/download/1/6/5/165255E7-1014-4D0A-B094-B6A430A6BFFC/vcredist_x64.exe' -OutFile $installer + $installation = Start-Process -FilePath $installer -ArgumentList '/q', '/norestart' -Wait -PassThru -WindowStyle Hidden + if ($installation.ExitCode -notin @(0, 3010) -or -not (Test-Path -LiteralPath $runtimeDll)) { + throw "VC++ 2010 x64 runtime installation failed (exit code $($installation.ExitCode))." + } + } + + & (Join-Path $hadoopBin 'winutils.exe') ls $hadoopBin + if ($LASTEXITCODE -ne 0) { + throw "Hadoop Windows tools failed to start (exit code $LASTEXITCODE)." + } + displayName: '[ENV] Prepare Hadoop for Windows Scala tests' + - task: Maven@3 displayName: 'Maven build Spark 4.0 bridge (JDK 17)' inputs: @@ -174,6 +202,9 @@ extends: javaHomeOption: 'JDKVersion' jdkVersionOption: '1.17' jdkArchitectureOption: 'x64' + env: + HADOOP_HOME: $(Build.BinariesDirectory)\hadoop + PATH: $(Build.BinariesDirectory)\hadoop\bin;$(PATH) - task: Maven@3 displayName: 'Maven build benchmark' From 67aa9c0da08deb674146051d0065921e24c36883 Mon Sep 17 00:00:00 2001 From: Shinai Yang Date: Thu, 17 Sep 2026 12:31:37 +0800 Subject: [PATCH 3/4] docs: clarify Windows prerequisites for Spark 4 tests --- docs/building/windows-instructions.md | 21 +++++++++++---------- 1 file changed, 11 insertions(+), 10 deletions(-) diff --git a/docs/building/windows-instructions.md b/docs/building/windows-instructions.md index 53b744153..1f1a58e50 100644 --- a/docs/building/windows-instructions.md +++ b/docs/building/windows-instructions.md @@ -61,19 +61,18 @@ If you already have all the pre-requisites, skip to the [build](windows-instruct - 6. Install **[WinUtils](https://github.com/steveloughran/winutils)** - - Download `winutils.exe` binary from [WinUtils repository](https://github.com/steveloughran/winutils). Select the Hadoop version bundled with the chosen Spark 3.5.3 distribution; do not reuse Hadoop 2.7 binaries from an older Spark installation. - - Save `winutils.exe` binary to a directory of your choice e.g., `c:\hadoop\bin` - - Set `HADOOP_HOME` to reflect the directory with winutils.exe (without bin). For instance, using command-line: - ```powershell - set HADOOP_HOME=c:\hadoop - ``` - - Set PATH environment variable to include `%HADOOP_HOME%\bin`. For instance, using command-line: + 6. Install **Hadoop Windows tools (WinUtils)** + - Use Hadoop Windows binaries compatible with your Spark distribution. To use the same binaries as Windows CI, download the [Hadoop 3.3.5 tools](https://github.com/SparkSnail/winutils/releases/download/hadoop-3.3.5/hadoop-3.3.5.zip) and copy `winutils.exe` and `hadoop.dll` from the archive's `hadoop-3.3.5` folder into a local directory, e.g., `C:\hadoop\bin`. Do not reuse Hadoop 2.7 binaries from an older Spark installation. + - These CI binaries require the **Microsoft Visual C++ 2010 x64 Redistributable** (`MSVCR100.dll`). If that runtime is missing, install it using the [Microsoft installer](https://download.microsoft.com/download/1/6/5/165255E7-1014-4D0A-B094-B6A430A6BFFC/vcredist_x64.exe) before running the tools. + - Set `HADOOP_HOME` to the directory containing `bin`, add its `bin` directory to `PATH`, and verify that `winutils.exe` starts in the current PowerShell session: + ```powershell - set PATH=%HADOOP_HOME%\bin;%PATH% + $env:HADOOP_HOME = 'C:\hadoop' + $env:PATH = "$env:HADOOP_HOME\bin;$env:PATH" + & "$env:HADOOP_HOME\bin\winutils.exe" ls "$env:HADOOP_HOME\bin" + if ($LASTEXITCODE -ne 0) { throw 'Hadoop Windows tools failed to start.' } ``` - Please make sure you are able to run `dotnet`, `java`, `mvn`, `spark-shell` from your command-line before you move to the next section. Feel there is a better way? Please [open an issue](https://github.com/dotnet/spark/issues) and feel free to contribute. > **Note**: A new instance of the command-line may be required if any environment variables were updated. @@ -104,6 +103,8 @@ Spark 4.0 is built separately. From the repository root, switch `JAVA_HOME` and mvn -f src/scala/microsoft-spark-4-0/pom.xml clean package ``` +`clean package` includes Spark 4 tests that use `SparkContext.addFile`. On Windows, complete the WinUtils setup in prerequisite 6 before running this command. + Its bridge is `src\scala\microsoft-spark-4-0\target\microsoft-spark-4-0_2.13-.jar` and requires a Scala 2.13 Spark 4.0 runtime with JDK 17. Building this bridge does not prove full API or deployment compatibility. See the [migration guide](../migration-guide.md#upgrading-after-spark-2x-support-removal) for package matching and clean-output requirements. ## Building .NET Samples Application From 2c30ccab6e573b569f200f52d95a421f880933be Mon Sep 17 00:00:00 2001 From: Shinai Yang Date: Thu, 17 Sep 2026 16:54:34 +0800 Subject: [PATCH 4/4] Fix DataFrameWriterV2 overwrite predicate in E2E tests --- .../IpcTests/Sql/DataFrameWriterV2Tests.cs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterV2Tests.cs b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterV2Tests.cs index 7f8ac7aa5..d6552d71d 100644 --- a/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterV2Tests.cs +++ b/src/csharp/Microsoft.Spark.E2ETest/IpcTests/Sql/DataFrameWriterV2Tests.cs @@ -55,7 +55,8 @@ public void TestSignaturesV3_0_X() AssertUnsupportedTableOperation(() => dfwV2.Replace()); AssertUnsupportedTableOperation(() => dfwV2.CreateOrReplace()); AssertUnsupportedTableOperation(() => dfwV2.Append()); - AssertUnsupportedTableOperation(() => dfwV2.Overwrite(df.Col("age"))); + // Use an unbound false predicate to reach the overwrite-by-filter capability check. + AssertUnsupportedTableOperation(() => dfwV2.Overwrite(Functions.Lit(false))); AssertUnsupportedTableOperation(() => dfwV2.OverwritePartitions()); }