Zum Hauptinhalt springen

Schreiben Sie Ihren ersten Mesh Adapter

Dieses Tutorial hilft Ihnen, Ihren ersten Mesh Adapter für OctoMesh zu schreiben.

Der Kürze halber gehen wir davon aus, dass Sie Ihre Entwicklungsumgebung bereits eingerichtet haben und über ein Grundverständnis von C# und .NET Core verfügen.

info

Wir gehen außerdem davon aus, dass Sie ein C#-Projekt eingerichtet haben, wie im Abschnitt Adapter-Projekt erstellen dokumentiert.

hinweis

Eine Referenzimplementierung finden Sie im folgenden Repository

Überblick​

Wir führen die folgenden Schritte aus:

  1. Konfiguration zum Ausführen des Adapters hinzufügen
  2. Adapter-Logik implementieren
  3. Pipeline-Logik implementieren

Mesh Adapter sind ausschließlich für die Ausführung in der Cloud konzipiert. Sie können sich direkt mit den OctoMesh-Repositories (MongoDB + CRATE.IO) verbinden. Sie können Daten von einem Edge Adapter oder anderen Quellen empfangen.

1. Konfiguration zum Ausführen des Adapters hinzufügen​

Navigieren Sie zur AdminUI und erstellen Sie einen neuen Adapter:

Einen neuen Mesh Adapter erstellen

Nach dem Erstellen des Adapters benötigen wir den Runtime-Identifier (rtId), um den Adapter im Code zu konfigurieren.

RtId kopieren

Mit diesen Informationen in der Zwischenablage können wir zum nächsten Teil übergehen.

2. Adapter-Logik implementieren​

Aktualisieren Sie das Launch-Profil, um die Adapter-Konfiguration einzuschließen.

{
"$schema": "http://json.schemastore.org/launchsettings.json",
"profiles": {
"AdapterEdgeDemo": {
"commandName": "Project",
"environmentVariables": {
"ASPNETCORE_ENVIRONMENT": "Development",
"OCTO_ADAPTER__TENANTID": meshtest,
"OCTO_ADAPTER__ADAPTERRTID": 6760711ec4ff02221e0b532e,
"OCTO_ADAPTER__ADAPTERCKTYPEID": "System.Communication/Adapter",
"OCTO_SYSTEM__AdminUserPassword": "OctoAdmin1",
"OCTO_SYSTEM__DatabaseUserPassword": "OctoUser1",
"OCTO_SYSTEM__UseDirectConnection": true
}
}
}
}

Das OCTO_ADAPTER__TENANTID ist die Tenant-ID des Adapters. Das OCTO_ADAPTER__ADAPTERRTID ist der Runtime-Identifier des Adapters, der zur Identifizierung des Adapters im System verwendet wird. Ersetzen Sie die Werte durch die aus der AdminUI kopierten.

Das OCTO_ADAPTER__ADAPTERCKTYPEID ist der Typ des Adapters, nämlich System.Communication/Adapter.

Das SYSTEM__AdminUserPassword und das SYSTEM__DatabaseUserPassword sind die Passwörter für den Admin- und den Datenbankbenutzer. Da wir den Adapter im Entwicklungsmodus ausführen, können wir die Standardpasswörter verwenden.

info

Mesh Adapter verbinden sich aus Performance-Gründen direkt mit den OctoMesh-Repositories.

Beginnen Sie in Program.cs mit der Konfiguration des Adapters.

// WebAdapterBuilder is a builder for creating Adapters acting as a Socket (Listener) or a Web API (Host)
var adapterBuilder = new WebAdapterBuilder();

await adapterBuilder.RunAsync(args, builder =>
{
// Define the configuration for the adapter
builder.Services.Configure<OctoSystemConfiguration>(options =>
builder.Configuration.GetSection("System").Bind(options));

// Add services to the container.

// Add the adapter service to startup and shutdown the adapter
builder.Services.AddSingleton<IAdapterService, AdapterMeshDemoService>();

// Add mesh adapter nodes and services to the container
builder.Services.AddOctoMeshAdapter();

}, app =>
{
app.UseOctoMeshAdapter();
});

Als Nächstes müssen wir das Interface IAdapterService implementieren, eine neue Klasse AdapterMeshDemoService erstellen und das Interface IAdapterService wie unten gezeigt implementieren.

using Meshmakers.Octo.Common.DistributionEventHub.Services;
using Meshmakers.Octo.Communication.Contracts.DataTransferObjects;
using Meshmakers.Octo.Sdk.Common.Adapters;
using Meshmakers.Octo.Sdk.Common.Services;
using Microsoft.Extensions.Logging;

namespace Meshmakers.Octo.Communication.MeshAdapter.Demo.Services;

internal class AdapterMeshDemoService(
ILogger<AdapterMeshDemoService> logger,
IPipelineRegistryService pipelineRegistryService,
IEventHubControl eventHubControl) : IAdapterService
{
public Task<bool> StartupAsync(AdapterStartup adapterStartup, List<DeploymentUpdateErrorMessageDto> errorMessages,
CancellationToken stoppingToken)
{
logger.LogInformation("Startup of mesh adapter");
try
{
return Task.FromResult(true);
}
catch (Exception e)
{
logger.LogError(e, "Error while startup");
throw;
}
}

public async Task ShutdownAsync(AdapterShutdown adapterShutdown, CancellationToken stoppingToken)
{
try
{
logger.LogInformation("Shutdown of mesh adapter");


logger.LogInformation("Mesh Adapter service stopped");
}
catch (Exception e)
{
logger.LogError(e, "Error while shutdown");
throw;
}
}
}

Nach diesen Schritten können Sie den Adapter ausführen und die Log-Ausgabe sehen.

In der Admin UI sollten Sie sehen, dass der Adapter online ist!

Adapter online

Nun fahren wir mit dem nächsten Schritt fort.

3. Node-Logik implementieren​

Zuerst müssen wir ein Node-Konfigurations-Record erstellen. Dieser Record wird verwendet, um Konfigurationsdaten für den Node zu serialisieren und zu deserialisieren.

[NodeName("Demo", 1)]
public record DemoNodeConfiguration : SourceTargetPathNodeConfiguration
{
}

NodeName definiert den Namen und die Version des Nodes. Ein Adapter kann mehrere Nodes mit demselben Namen, aber unterschiedlichen Versionen haben. Typischerweise hat ein Node Eingabe- und Ausgabedaten. Das SourceTargetPathNodeConfiguration wird verwendet, um den Eingabe- und Ausgabedatenpfad (Path, TargetPath) für den Node zu definieren.

Extract- und Load-Nodes können nur Eingabedaten oder nur Ausgabedaten haben. In diesem Fall wird das PathNodeConfiguration oder das TargetPathNodeConfiguration verwendet.

Zweitens müssen wir den Node selbst erstellen.

[NodeConfiguration(typeof(DemoNodeConfiguration))]
public class DemoNode(NodeDelegate next) : IPipelineNode
{
public async Task ProcessObjectAsync(IDataContext dataContext, INodeContext nodeContext)
{
// Continue with next node in pipeline
await next(dataContext, nodeContext);
}
}

Das Attribut NodeConfiguration wird verwendet, um das Konfigurations-Record für den Node zu definieren.

Drittens müssen wir die Startup-Konfiguration anpassen, um den Node einzuschließen.


// WebAdapterBuilder is a builder for creating Adapters acting as a Socket (Listener) or a Web API (Host)
var adapterBuilder = new WebAdapterBuilder();

await adapterBuilder.RunAsync(args, builder =>
{
// Define the configuration for the adapter
builder.Services.Configure<OctoSystemConfiguration>(options =>
builder.Configuration.GetSection("System").Bind(options));

builder.Services.Configure<MeshAdapterConfiguration>(options =>
builder.Configuration.GetSection("Adapter").Bind(options));

// Add services to the container.

// Add the adapter service to startup and shutdown the adapter
builder.Services.AddSingleton<IAdapterService, AdapterMeshDemoService>();

// Add mesh adapter nodes and services to the container
builder.Services.AddOctoMeshAdapter()
.RegisterNode<DemoNode>(); // Sample to register a node

}, app =>
{
app.UseOctoMeshAdapter();
});

Als Nächstes müssen wir die Node-Logik implementieren. Wir erweitern die Node-Konfiguration um eine neue Property MyMessage.

[NodeName("Demo", 1)]
public record DemoNodeConfiguration : TargetPathNodeConfiguration
{
public required string MyMessage { get; set; } = "Hello, World!";
}

Wir möchten die Nachricht in den Zielpfad schreiben.

[NodeConfiguration(typeof(DemoNodeConfiguration))]
public class DemoNode(NodeDelegate next) : IPipelineNode
{
public async Task ProcessObjectAsync(IDataContext dataContext, INodeContext nodeContext)
{
// Get configuration
var c = nodeContext.GetNodeConfiguration<DemoNodeConfiguration>();

// set value
dataContext.SetValueByPath(c.TargetPath, c.DocumentMode, c.TargetValueKind, c.TargetValueWriteMode,
c.MyMessage);

// Continue with next node in pipeline
await next(dataContext, nodeContext);
}
}

Nun müssen wir den Adapter-Service aktualisieren, um Folgendes einzuschließen:

  • Registrierung von Pipelines, die ausgeführt werden müssen
  • Start von Trigger-Nodes, die die Pipeline bei einem bestimmten Event ausführen
internal class AdapterMeshDemoService(
ILogger<AdapterMeshDemoService> logger,
IPipelineRegistryService pipelineRegistryService,
IEventHubControl eventHubControl) : IAdapterService
{
public async Task<bool> StartupAsync(AdapterStartup adapterStartup, List<DeploymentUpdateErrorMessageDto> errorMessages,
CancellationToken stoppingToken)
{
logger.LogInformation("Startup of mesh adapter");
try
{
var success =await pipelineRegistryService.RegisterPipelinesAsync(adapterStartup.TenantId,
adapterStartup.Configuration.Pipelines, errorMessages);
await pipelineRegistryService.StartTriggerPipelineNodesAsync(adapterStartup.TenantId);

await eventHubControl.StartAsync(stoppingToken);

return success;
}
catch (Exception e)
{
logger.LogError(e, "Error while startup");
throw;
}
}

public async Task ShutdownAsync(AdapterShutdown adapterShutdown, CancellationToken stoppingToken)
{
try
{
logger.LogInformation("Shutdown of mesh adapter");
await pipelineRegistryService.StopTriggerPipelineNodesAsync(adapterShutdown.TenantId);

pipelineRegistryService.UnregisterAllPipelines(adapterShutdown.TenantId);
await eventHubControl.StopAsync(stoppingToken);
logger.LogInformation("Mesh Adapter service stopped");
}
catch (Exception e)
{
logger.LogError(e, "Error while shutdown");
throw;
}
}
}

Ihr Adapter ist nun bereit, die Pipeline auszuführen. Der letzte Schritt des Puzzles besteht darin, die Datenpipeline in der AdminUI zu konfigurieren. Wir müssen eine neue Datenpipeline erstellen.

Die Pipeline erstellen

Ein DataFlow gruppiert zusammengehörige Pipelines, die als Teil eines Datenverarbeitungs-Workflows zusammenarbeiten. Pipelines können je nach Anforderung auf Edge Adaptern oder Cloud Adaptern ausgeführt werden.

Wir haben dem Adapter nun mitgeteilt, die Pipeline auszuführen, aber wir müssen die Pipeline noch konfigurieren, was ebenfalls in der AdminUI erfolgen kann.

Wenn Sie nach dem Speichern der Datenpipeline auf die Schaltfläche Edit/Debug definitions klicken, sehen Sie die Pipeline Builder UI, in der Sie die Pipeline-Definition eingeben.

Eine einfache Pipeline-Definition könnte so aussehen:

triggers:
- type: FromExecutePipelineCommand@1
transformations:
- type: Demo@1
description: Simulates data
myMessage: Hello, Mars!

Diese Pipeline-Definition enthält einen Trigger-Node, der es ermöglicht, die Pipeline aus der AdminUI heraus auszuführen, sowie einen Transformations-Node, der die Nachricht „Hello, Mars!" in den Zielpfad schreibt.

Speichern Sie die Pipeline und stellen Sie die Pipeline bereit (Download) Definition auf dem Adapter.

Die Pipeline bereitstellen

Nun können Sie den Adapter starten und die Pipeline aus der AdminUI heraus ausführen.

Die Pipeline ausführen

Sie sollten die Nachricht in der Log-Ausgabe des Adapters sehen. Dies ist ein Auszug aus der Log-Ausgabe:

INFO|FromExecutePipelineCommand@1: Received command executing pipeline
DEBUG|[meshtest] Running pipeline for pipeline System.Communication/Pipeline@67e135a9477e78e980bbb516 as run with execution id 7a47b1cb-87ae-4511-8ba5-4a87e5825987
INFO|PipelineExecution: Executing pipeline
DEBUG|PipelineExecution: Node completed
DEBUG|PipelineExecution/Demo@1: Forward Executing
DEBUG|PipelineExecution/Demo@1: Node completed
DEBUG|PipelineExecution/Demo@1: Reverse completed
INFO|PipelineExecution: Pipeline completed
DEBUG|PipelineExecution: Node completed
warnung

Eine Referenzimplementierung finden Sie im folgenden Repository. Dieses Repository enthält eine fortgeschrittenere Adapter-Implementierung mit Containern, einschließlich einer OctoMesh-ETL-Pipeline und einer CI/CD-Pipeline auf Basis von Azure DevOps