Zum Hauptinhalt springen

Schreiben Sie Ihren ersten Edge Adapter

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

Der Kürze halber gehen wir davon aus, dass Sie Ihre Entwicklungsumgebung bereits eingerichtet haben und über ein grundlegendes Verstä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 Creating Adapter Project dokumentiert.

hinweis

Für eine Referenzimplementierung können Sie das folgende Repository prüfen

Ü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

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

Navigieren Sie zur AdminUI und erstellen Sie einen neuen Adapter:

Create a new Edge Adapter

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

Copy RtId

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

2. Adapter-Logik implementieren​

Aktualisieren Sie das Launch-Profil, um die Adapter-Konfiguration aufzunehmen.

{
"$schema": "http://json.schemastore.org/launchsettings.json",
"profiles": {
"AdapterEdgeDemo": {
"commandName": "Project",
"environmentVariables": {
"ASPNETCORE_ENVIRONMENT": "Development",
"OCTO_ADAPTER__TENANTID": meshtest,
"OCTO_ADAPTER__ADAPTERRTID": 6760711ec4ff02221e0b532d
}
}
}
}

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

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

using Meshmakers.Octo.Communication.EdgeAdapter.Demo.Services;
using Meshmakers.Octo.Sdk.Common.Adapters;
using Meshmakers.Octo.Sdk.Common.EtlDataPipeline;
using Meshmakers.Octo.Sdk.Common.Services;

// AdapterBuilder is used for active adapters (connecting to external systems)
var adapterBuilder = new AdapterBuilder();

adapterBuilder.Run(args, (_, services) =>
{
// Register data pipeline services
services.AddDataPipeline()
// Register the ETL context to access meta data of the execution (e.g. tenant id, pipeline id, ...)
.RegisterEtlContext<IEtlContext>();
services.AddSingleton<IAdapterService, AdapterEdgeDemoService>();
});

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

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

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

/// <summary>
/// This is the main service for the plug. It is responsible for starting and stopping the plug.
/// Shutdown is called when the plug is stopped or new configuration is received.
/// Startup is called when the plug starts or receives a new configuration.
/// </summary>
/// <param name="pipelineRegistryService">Service for registering and starting pipelines</param>
/// <param name="eventHubControl">Event hub control service</param>
public class AdapterEdgeDemoService(
ILogger<AdapterEdgeDemoService> logger,
IPipelineRegistryService pipelineRegistryService,
IEventHubControl eventHubControl)
: IAdapterService
{
public Task<bool> StartupAsync(AdapterStartup adapterStartup, List<DeploymentUpdateErrorMessageDto> errorMessages, CancellationToken stoppingToken)
{
logger.LogInformation("Startup");

try
{
return Task.FromResult(true);
}
catch (Exception e)
{
logger.LogError(e, "Error while startup");
throw;
}
}

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

logger.LogInformation("Shutdown complete");
return Task.CompletedTask;
}
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 einen Node-Konfigurationsrecord 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. Die 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 die PathNodeConfiguration oder 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 den Konfigurationsrecord für den Node zu definieren.

Drittens müssen wir die Startup-Konfiguration ändern, um den Node aufzunehmen.

// AdapterBuilder is used for active adapters (connecting to external systems)
var adapterBuilder = new AdapterBuilder();

adapterBuilder.Run(args, (_, services) =>
{
// Register data pipeline services
services.AddDataPipeline()
.RegisterNode<DemoNode>() // Sample to register a node
// Register the ETL context to access meta data of the execution (e.g. tenant id, pipeline id, ...)
.RegisterEtlContext<IEtlContext>();
services.AddTransient<IPollingService, PollingService>();
services.AddSingleton<IAdapterService, AdapterEdgeDemoService>();
});

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 Target-Path 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 aufzunehmen:

  • Registrierung von Pipelines, die ausgeführt werden müssen
  • Start von Trigger-Nodes, die die Pipeline bei einem bestimmten Ereignis ausführen
public class AdapterEdgeDemoService(
ILogger<AdapterEdgeDemoService> logger,
IPipelineRegistryService pipelineRegistryService,
IEventHubControl eventHubControl)
: IAdapterService
{
public async Task<bool> StartupAsync(AdapterStartup adapterStartup, List<DeploymentUpdateErrorMessageDto> errorMessages, CancellationToken stoppingToken)
{
logger.LogInformation("Startup");

try
{
// adapterStartup contains configuration:
// Adapter configuration (optional) and a list of pipelines to be executed by this adapter.
// Pipelines are a sequence of nodes that process data.
// The pipeline is registered with the pipeline registry service.

// Register pipelines
var success = await pipelineRegistryService.RegisterPipelinesAsync(adapterStartup.TenantId,
adapterStartup.Configuration.Pipelines, errorMessages);

// If success is false, at least one pipeline failed to register, and the errorMessages list contains the error messages.
// The adapter should start to execute the rest of the pipelines.

// Start triggers. Triggers are special nodes that start the pipeline execution based on some event.
await pipelineRegistryService.StartTriggerPipelineNodesAsync(adapterStartup.TenantId);

// Start connection to rabbitmq event hub
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");

// Stop triggers
await pipelineRegistryService.StopTriggerPipelineNodesAsync(adapterShutdown.TenantId);

// Unregister pipelines
pipelineRegistryService.UnregisterAllPipelines(adapterShutdown.TenantId);

// Stop connection rabbitmq event hub
await eventHubControl.StopAsync(stoppingToken);

logger.LogInformation("Shutdown complete");
}
catch (Exception e)
{
logger.LogError(e, "Error while shutdown");
throw;
}
}
}

Ihr Adapter ist nun bereit, die Pipeline auszuführen. Der letzte Schritt im Puzzle besteht darin, die Datenpipeline in der AdminUI zu konfigurieren. Wir müssen einen neuen DataFlow erstellen.

Create the pipeline

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

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

Indem Sie nach dem Speichern der Datenpipeline den Button Edit/Debug definitions betätigen, sehen Sie die Pipeline-Builder-UI, in die Sie die Pipeline-Definition eingeben.

Eine grundlegende Pipeline-Definition könnte so aussehen:

triggers:
- type: FromPolling@1
interval: 00:00:10
transformations:
- type: Demo@1
description: Simulates data
myMessage: Hello, Mars!

Diese Pipeline-Definition enthält einen Trigger-Node, der alle 10 Sekunden pollt, und einen Transformations-Node, der die Nachricht "Hello, Mars!" in den Target-Path schreibt.

Speichern Sie die Pipeline und deployen (laden) Sie die Pipeline-Definition auf den Adapter.

Deploy the pipeline

Nach 10 Sekunden sollten Sie die Nachricht in der Log-Ausgabe des Adapters sehen.

DEBUG|[meshtest] Running pipeline for pipeline System.Communication/Pipeline@67e132d6477e78e980bbb512 as run with execution id e2d7e1fd-444f-40da-917c-a00afb6a71a2
DEBUG|Connected: guest@localhost:5672/ (address: amqp://localhost:5672, local: 61536)
DEBUG|Endpoint Ready: rabbitmq://localhost/MacBookProvonG_MeshmakersOcto_bus_65jyyyyc8i3jzmg1bdqsixornw?temporary=true
INFO|Bus started: rabbitmq://localhost/
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
warnung

Für eine Referenzimplementierung können Sie das folgende Repository prüfen. Dieses Repository enthält eine fortgeschrittenere Adapter-Implementierung unter Verwendung von Containern, einschließlich einer OctoMesh-ETL-Pipeline und einer CI/CD-Pipeline auf Basis von Azure DevOps