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
| Typ | Deployment | Anwendungsfälle |
|---|---|---|
| Edge Adapter | Nahe an den Datenquellen (K3s) | Sensoren, PLCs, lokale Datenbanken, Dateisysteme |
| Mesh Adapter | Cloud/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
| Stufe | Zweck | Beispiele |
|---|---|---|
| Trigger | Löst die Pipeline-Ausführung aus | Zeitplan, HTTP-Anfrage, Nachricht, Entitätsänderung |
| Extract | Ruft Daten aus der Quelle ab | Datenbank abfragen, Datei lesen, API aufrufen |
| Transform | Verarbeitet und mappt Daten | Filtern, Aggregieren, Formatkonvertierung |
| Load | Speichert Daten im Ziel | In 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
- Erstellen Sie ein neues .NET-Projekt:
dotnet new console -n MyAdapter
cd MyAdapter
- 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
PipelineTriggerist ein Kind des DataFlow (System/ParentChild) und verweist über die AssoziationSystem.Communication/Triggersauf 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), mitSystem/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
- Edge-Adapter erstellen: Ausführliche Anleitung für Edge-Adapter
- Mesh-Adapter erstellen: Ausführliche Anleitung für Mesh-Adapter
- Trigger-Nodes: Eigene Trigger erstellen
- API-Referenz: Vollständige Node-Referenz