Zum Hauptinhalt springen

Adapter-Entwicklung

Dieses Dokument bietet einen umfassenden Leitfaden zur Entwicklung von Adaptern für OctoMesh und behandelt sowohl Edge- als auch Mesh-Adapter.

Adapter-Überblick​

Adapter sind die Integrationsschicht von OctoMesh und verbinden externe Systeme mit der Data-Mesh-Plattform.

Adapter-Typen​

TypDeploymentAnwendungsfälle
Edge AdapterNahe an den Datenquellen (K3s)Sensoren, PLCs, lokale Datenbanken, Dateisysteme
Mesh AdapterCloud/Rechenzentrum (K8s)ERP-Systeme, Cloud-Dienste, Datenaggregation

Adapter-Rollen​

  • Plugs: Holen Daten IN OctoMesh hinein
  • Sockets: Stellen Daten AUS OctoMesh bereit

Architektur der Datenpipeline​

Adapter führen Datenpipelines aus, die definieren, wie Daten durch das System fließen.

Pipeline-Struktur​

StufeZweckBeispiele
TriggerLöst die Pipeline-Ausführung ausZeitplan, HTTP-Anfrage, Nachricht, Entitätsänderung
ExtractRuft Daten aus der Quelle abDatenbank abfragen, Datei lesen, API aufrufen
TransformVerarbeitet und mappt DatenFiltern, Aggregieren, Formatkonvertierung
LoadSpeichert Daten im ZielIn RT speichern, in den Stream schreiben, Benachrichtigung senden

Pipeline-Nodes​

Pipelines bestehen aus Nodes, die bestimmte Operationen ausführen:

Ein Adapter-Projekt erstellen​

Voraussetzungen​

  • .NET 8.0 SDK oder neuer
  • Visual Studio, JetBrains Rider oder VS Code
  • Zugriff auf die OctoMesh NuGet-Pakete

Projekteinrichtung​

  1. Erstellen Sie ein neues .NET-Projekt:
dotnet new console -n MyAdapter
cd MyAdapter
  1. Fügen Sie die erforderlichen NuGet-Pakete hinzu:
# Core SDK packages
dotnet add package Meshmakers.Octo.Sdk.Common
dotnet add package Meshmakers.Octo.Sdk.ServiceClient

# For Mesh Adapter
dotnet add package Meshmakers.Octo.Sdk.MeshAdapter

# For Edge Adapter
dotnet add package Meshmakers.Octo.Sdk.EdgeAdapter

Projektstruktur​

MyAdapter/
├── Program.cs # Entry point
├── appsettings.json # Configuration
├── Nodes/
│ ├── Trigger/ # Trigger nodes
│ ├── Extract/ # Extract nodes
│ ├── Transform/ # Transform nodes
│ └── Load/ # Load nodes
├── Pipelines/
│ └── MyPipeline.cs # Pipeline definitions
└── Services/
└── MyService.cs # Custom services

Pipeline-Nodes implementieren​

Beispiel für einen Trigger-Node​

Ein Trigger-Node löst die Pipeline-Ausführung aus:

using Meshmakers.Octo.Sdk.Common.Pipelines;

public class FromHttpRequestNodeConfiguration : ITriggerNodeConfiguration
{
public string Route { get; set; } = "/api/trigger";
public HttpMethod Method { get; set; } = HttpMethod.Post;
public bool RequireAuthentication { get; set; } = true;
}

public class FromHttpRequestNode : TriggerNode<FromHttpRequestNodeConfiguration>
{
public override async Task<PipelineContext> ExecuteAsync(
FromHttpRequestNodeConfiguration config,
CancellationToken cancellationToken)
{
// Access HTTP request data
var requestBody = Context.GetHttpRequestBody();

// Create pipeline context with data
var context = new PipelineContext();
context.SetData("request", requestBody);

return context;
}
}

Beispiel für einen Extract-Node​

Ein Extract-Node ruft Daten ab:

public class GetRtEntitiesByTypeNodeConfiguration : IExtractNodeConfiguration
{
public string CkTypeId { get; set; }
public int MaxResults { get; set; } = 100;
public string Filter { get; set; }
}

public class GetRtEntitiesByTypeNode : ExtractNode<GetRtEntitiesByTypeNodeConfiguration>
{
private readonly IAssetRepositoryClient _client;

public GetRtEntitiesByTypeNode(IAssetRepositoryClient client)
{
_client = client;
}

public override async Task<PipelineContext> ExecuteAsync(
GetRtEntitiesByTypeNodeConfiguration config,
PipelineContext context,
CancellationToken cancellationToken)
{
var entities = await _client.QueryEntitiesAsync(
config.CkTypeId,
config.Filter,
config.MaxResults,
cancellationToken);

context.SetData("entities", entities);
return context;
}
}

Beispiel für einen Transform-Node​

Ein Transform-Node verarbeitet Daten:

public class DataMappingNodeConfiguration : ITransformNodeConfiguration
{
public List<MappingEntry> Mappings { get; set; } = new();
}

public class MappingEntry
{
public string SourcePath { get; set; }
public string TargetPath { get; set; }
public string TransformExpression { get; set; }
}

public class DataMappingNode : TransformNode<DataMappingNodeConfiguration>
{
public override async Task<PipelineContext> ExecuteAsync(
DataMappingNodeConfiguration config,
PipelineContext context,
CancellationToken cancellationToken)
{
var sourceData = context.GetData<object>("source");
var result = new Dictionary<string, object>();

foreach (var mapping in config.Mappings)
{
var value = ExtractValue(sourceData, mapping.SourcePath);

if (!string.IsNullOrEmpty(mapping.TransformExpression))
{
value = ApplyTransform(value, mapping.TransformExpression);
}

SetValue(result, mapping.TargetPath, value);
}

context.SetData("mapped", result);
return context;
}
}

Beispiel für einen Load-Node​

Ein Load-Node speichert Daten:

public class ApplyChangesNodeConfiguration : ILoadNodeConfiguration
{
public UpdateKind UpdateKind { get; set; } = UpdateKind.CreateOrUpdate;
public bool ValidateBeforeSave { get; set; } = true;
}

public class ApplyChangesNode : LoadNode<ApplyChangesNodeConfiguration>
{
private readonly IAssetRepositoryClient _client;

public ApplyChangesNode(IAssetRepositoryClient client)
{
_client = client;
}

public override async Task<PipelineContext> ExecuteAsync(
ApplyChangesNodeConfiguration config,
PipelineContext context,
CancellationToken cancellationToken)
{
var updateInfo = context.GetData<UpdateInfo>("updateInfo");

var result = await _client.ApplyChangesAsync(
updateInfo,
config.UpdateKind,
cancellationToken);

context.SetData("result", result);
return context;
}
}

Pipeline-Konfiguration​

Eine Pipeline definieren​

Pipelines werden als YAML verfasst und auf einer System.Communication/Pipeline-Runtime-Entität im Attribut System.Communication/PipelineDefinition gespeichert. Es gibt kein JSON-Pipeline-Format und keine C#-API zur Pipeline-Konstruktion — der Adapter liest die YAML-Definition, wenn die Pipeline bereitgestellt wird. Die verfügbaren Node-Typen finden Sie in der Referenz der Pipeline-Nodes.

triggers:
- type: FromPipelineTriggerEvent@1
transformations:
- type: Simulation@1
description: Generate random sensor values
locale: en
simulations:
- targetPath: $.sensor.temperature
simulatorKey: Math.DoubleRandom
configuration: "{min:18, max:32}"
- type: Logger@1
description: Log generated data
message: Scheduled trigger executed - sensor data generated

Eine Pipeline planen​

Es gibt keinen FromSchedule-Trigger-Node. Cron-basierte Zeitplanung wird außerhalb der Pipeline-Definition modelliert, auf einer separaten System.Communication/PipelineTrigger-Entität:

  • Der PipelineTrigger ist ein Kind des DataFlow (System/ParentChild) und verweist über die Assoziation System.Communication/Triggers auf die Pipelines, die er auslöst.
  • Der Cron-Ausdruck liegt im Attribut System.Bot/CronExpression (6-Feld-Quartz-Stil, z. B. "0 * * * * ?" für jede Minute), mit System/Enabled, um ihn einzuschalten.
  • Der Communication Controller registriert den wiederkehrenden Job und veröffentlicht planmäßig ein Trigger-Event; die Pipeline empfängt es über einen FromPipelineTriggerEvent@1-Trigger-Node in ihrer Definition.
# PipelineTrigger runtime entity (sibling of the Pipeline, child of the DataFlow)
- rtId: dd0000000000000000000003
ckTypeId: System.Communication/PipelineTrigger
associations:
- roleId: System/ParentChild
targetRtId: dd0000000000000000000001 # the DataFlow
targetCkTypeId: System.Communication/DataFlow
- roleId: System.Communication/Triggers
targetRtId: dd0000000000000000000002 # the Pipeline to fire
targetCkTypeId: System.Communication/Pipeline
attributes:
- id: System/Enabled
value: true
- id: System.Bot/CronExpression
value: "0 * * * * ?"

Für adapterlokale periodische Ausführung ohne zentralen Zeitplan verwenden Sie stattdessen FromPolling@1 — dessen interval ist eine TimeSpan (00:05:00 für alle fünf Minuten, 1.00:00:00 für täglich), kein Cron-Ausdruck.

Adapter-Konfiguration​

appsettings.json​

{
"OctoMesh": {
"IdentityServiceUrl": "https://identity.octomesh.local",
"AssetRepositoryUrl": "https://asset-repo.octomesh.local",
"CommunicationControllerUrl": "https://com-controller.octomesh.local",
"TenantId": "my-tenant",
"ClientId": "my-adapter",
"ClientSecret": "${secrets.clientSecret}"
},
"Adapter": {
"PoolId": "edge-pool-01",
"AdapterId": "energy-adapter-01"
},
"Logging": {
"LogLevel": {
"Default": "Information",
"Meshmakers": "Debug"
}
}
}

Einrichtung der Dependency Injection​

public class Program
{
public static async Task Main(string[] args)
{
var builder = Host.CreateApplicationBuilder(args);

// Configure OctoMesh services
builder.Services.AddOctoMeshClient(options =>
{
options.IdentityServiceUrl = builder.Configuration["OctoMesh:IdentityServiceUrl"];
options.AssetRepositoryUrl = builder.Configuration["OctoMesh:AssetRepositoryUrl"];
options.TenantId = builder.Configuration["OctoMesh:TenantId"];
});

// Register adapter
builder.Services.AddMeshAdapter(options =>
{
options.PoolId = builder.Configuration["Adapter:PoolId"];
options.AdapterId = builder.Configuration["Adapter:AdapterId"];
});

// Register custom nodes
builder.Services.AddTransient<MyCustomNode>();

var host = builder.Build();
await host.RunAsync();
}
}

Logging in Nodes​

Implementieren Sie ordnungsgemäßes Logging für Debugging und Monitoring:

public class MyCustomNode : TransformNode<MyCustomNodeConfiguration>
{
private readonly ILogger<MyCustomNode> _logger;

public MyCustomNode(ILogger<MyCustomNode> logger)
{
_logger = logger;
}

public override async Task<PipelineContext> ExecuteAsync(
MyCustomNodeConfiguration config,
PipelineContext context,
CancellationToken cancellationToken)
{
_logger.LogInformation("Starting transformation with config: {@Config}", config);

try
{
// Node logic here
var result = await ProcessDataAsync(context, cancellationToken);

_logger.LogInformation("Transformation completed. Processed {Count} items", result.Count);

context.SetData("result", result);
return context;
}
catch (Exception ex)
{
_logger.LogError(ex, "Transformation failed");
throw;
}
}
}

Fehlerbehandlung​

Pipeline-Fehlerbehandlung​

public class MyLoadNode : LoadNode<MyLoadNodeConfiguration>
{
public override async Task<PipelineContext> ExecuteAsync(
MyLoadNodeConfiguration config,
PipelineContext context,
CancellationToken cancellationToken)
{
try
{
await SaveDataAsync(context, cancellationToken);
}
catch (ValidationException ex)
{
// Mark as partial failure, continue pipeline
context.AddWarning($"Validation failed: {ex.Message}");
}
catch (ConnectionException ex)
{
// Retry logic
for (int i = 0; i < 3; i++)
{
await Task.Delay(TimeSpan.FromSeconds(Math.Pow(2, i)), cancellationToken);
try
{
await SaveDataAsync(context, cancellationToken);
break;
}
catch { /* Continue retry */ }
}
}

return context;
}
}

Adapter testen​

Unit-Tests für Nodes​

[TestClass]
public class DataMappingNodeTests
{
[TestMethod]
public async Task ExecuteAsync_WithValidMapping_ShouldMapData()
{
// Arrange
var node = new DataMappingNode();
var config = new DataMappingNodeConfiguration
{
Mappings = new List<MappingEntry>
{
new() { SourcePath = "$.name", TargetPath = "displayName" }
}
};
var context = new PipelineContext();
context.SetData("source", new { name = "Test" });

// Act
var result = await node.ExecuteAsync(config, context, CancellationToken.None);

// Assert
var mapped = result.GetData<Dictionary<string, object>>("mapped");
Assert.AreEqual("Test", mapped["displayName"]);
}
}

Integrationstests​

[TestClass]
public class EnergyAdapterIntegrationTests
{
private IHost _host;

[TestInitialize]
public void Setup()
{
_host = Host.CreateDefaultBuilder()
.ConfigureServices((context, services) =>
{
services.AddOctoMeshClient(options =>
{
options.IdentityServiceUrl = "https://localhost:5003";
// ... test configuration
});
})
.Build();
}

[TestMethod]
public async Task Pipeline_ShouldSyncEnergyMeters()
{
// Test implementation
}
}

Deployment​

Docker-Container​

FROM mcr.microsoft.com/dotnet/aspnet:8.0 AS base
WORKDIR /app

FROM mcr.microsoft.com/dotnet/sdk:8.0 AS build
WORKDIR /src
COPY ["MyAdapter.csproj", "."]
RUN dotnet restore
COPY . .
RUN dotnet build -c Release -o /app/build

FROM build AS publish
RUN dotnet publish -c Release -o /app/publish

FROM base AS final
WORKDIR /app
COPY --from=publish /app/publish .
ENTRYPOINT ["dotnet", "MyAdapter.dll"]

Kubernetes-Deployment​

apiVersion: apps/v1
kind: Deployment
metadata:
name: energy-adapter
namespace: octo-adapters
spec:
replicas: 1
selector:
matchLabels:
app: energy-adapter
template:
metadata:
labels:
app: energy-adapter
spec:
containers:
- name: adapter
image: myregistry/energy-adapter:latest
env:
- name: OctoMesh__TenantId
value: "production"
- name: OctoMesh__ClientSecret
valueFrom:
secretKeyRef:
name: adapter-secrets
key: client-secret

Nächste Schritte​