Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 28 additions & 5 deletions azure-pipelines-e2e-tests-template.yml
Original file line number Diff line number Diff line change
Expand Up @@ -268,12 +268,14 @@ stages:
displayName: 'E2E tests'
inputs:
command: test
# Spark 4 covers the bridge, row retrieval, UDF, RDD, Broadcast, Worker lifecycle and ML. 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:
Expand Down Expand Up @@ -317,6 +319,23 @@ stages:
failTaskOnMissingResultsFile: true
testRunTitle: 'ML persistence Spark ${{ test.version }} (${{ option.pool }})'

- 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')) }}:
Expand Down Expand Up @@ -380,7 +399,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:
Expand All @@ -407,7 +428,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:
Expand Down
33 changes: 32 additions & 1 deletion azure-pipelines-pr.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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'
Expand Down Expand Up @@ -226,7 +257,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|Category=ML"'
testOptions: '--filter "Category=Spark40Compatibility|Category=DataFrameRowRetrieval|Category=Broadcast|Category=ArrowUdf|Category=Avro|Category=Streaming|Category=ML|Category=DriverApi"'
${{ else }}:
testOptions: ""
backwardCompatibleTestOptions: $(backwardCompatibleTestOptions_${{ pool }}_${{ split(version, '.')[0] }}_${{ split(version, '.')[1] }})
Expand Down
21 changes: 11 additions & 10 deletions docs/building/windows-instructions.md
Original file line number Diff line number Diff line change
Expand Up @@ -61,19 +61,18 @@ If you already have all the pre-requisites, skip to the [build](windows-instruct

</details>

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.
Expand Down Expand Up @@ -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-<version>.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
Expand Down
Original file line number Diff line number Diff line change
@@ -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<long, long>(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<long>(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<long, long>(LaterCell.Apply);
var fail = Udf<long, long>(LaterCell.Fail);
");
await SubmitAsync(kernel, @"
var later = spark.Range(0, 3, 1, 2).Select(multiply(Col(""id"")))
.Collect().Select(row => row.GetAs<long>(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<long, long>(FailedCell.Apply);
spark.Range(1).Select(fail(Col(""id""))).Collect().ToArray();
", "csharp"));
CommandFailed workerFailure = Assert.Single(failedAction.Events.OfType<CommandFailed>());
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<long>(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<CommandFailed>());
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<long>(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<CommandFailed>().Select(failure => failure.Message).ToArray();
Assert.True(failures.Length == 0, string.Join(Environment.NewLine, failures));
Assert.Single(result.Events.OfType<CommandSucceeded>());
}

private static void AssertValues(CSharpKernel kernel, string name, params long[] expected)
{
Assert.True(kernel.TryGetValue(name, out long[] values));
Assert.Equal(expected, values);
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
<Project Sdk="Microsoft.NET.Sdk">
<PropertyGroup>
<TargetFramework>net8.0</TargetFramework>
</PropertyGroup>
<ItemGroup>
<PackageReference Include="Microsoft.DotNet.Interactive.CSharp" Version="1.0.0-beta.25323.1" />
<ProjectReference Include="..\Microsoft.Spark.Extensions.DotNet.Interactive\Microsoft.Spark.Extensions.DotNet.Interactive.csproj" />
<ProjectReference Include="..\..\Microsoft.Spark.E2ETest\Microsoft.Spark.E2ETest.csproj" />
</ItemGroup>
</Project>
Loading
Loading