Quick Start Guide
Install the Go SDK and configure the [R]DP instance to use:
go mod init rdp-quickstart
go get github.com/raft-tech/[email protected]
export RDP_SERVER_URL=https://rdp-example.com
export RDP_API_KEY=my-rdp-api-key
Create a Transformer
package main
import (
"context"
"fmt"
"log"
rdpsdk "github.com/raft-tech/rdp-sdk-go"
transformersapi "github.com/raft-tech/rdp-sdk-go/transformers"
)
func main() {
cfg, err := rdpsdk.LoadConfig()
if err != nil {
log.Fatal(err)
}
client, err := rdpsdk.NewFromConfig(cfg)
if err != nil {
log.Fatal(err)
}
inputName := "Incoming messages"
inputType := "INTERNAL_STREAMING"
outputName := "Outgoing files"
outputType := "INTERNAL_OBJECT_STORE"
transformer := transformersapi.Transformer{
Uid: "new-transformer-abc",
Name: "My New Transformer",
Description: "Transforms data.",
Status: "available",
SecurityMarkings: "UNCLASSIFIED",
Types: []transformersapi.TransformerType{},
Inputs: map[string]transformersapi.DataConnTpl{
"INPUT_STREAMING_TOPIC": {
DisplayName: &inputName,
ConnType: &inputType,
},
},
Outputs: map[string]transformersapi.DataConnTpl{
"OUTPUT_OBJECT_BUCKET": {
DisplayName: &outputName,
ConnType: &outputType,
},
},
Configuration: transformersapi.ConfigurationTpl{
"environment": []map[string]any{
{
"name": "LOG_LEVEL",
"required": true,
"description": "The log level used at startup.",
},
},
},
Instantiation: transformersapi.CreationConfig{
"job_image": map[string]any{
"image": "my-container-image",
"pull_policy": "IfNotPresent",
"default_replicas": 1,
},
},
}
resp, err := client.Transformers().Create(context.Background(), transformer)
if err != nil {
log.Fatal(err)
}
if resp.JSON201 == nil {
log.Fatalf("create transformer failed: %s", resp.Status())
}
fmt.Printf("Created Transformer with UID %q\n", resp.JSON201.Uid)
}
Retrieve a Transformer
package main
import (
"context"
"fmt"
"log"
rdpsdk "github.com/raft-tech/rdp-sdk-go"
)
func main() {
cfg, err := rdpsdk.LoadConfig()
if err != nil {
log.Fatal(err)
}
client, err := rdpsdk.NewFromConfig(cfg)
if err != nil {
log.Fatal(err)
}
resp, err := client.Transformers().ID("new-transformer-abc").Get(context.Background())
if err != nil {
log.Fatal(err)
}
if resp.JSON200 == nil {
log.Fatalf("get transformer failed: %s", resp.Status())
}
fmt.Printf("Found transformer: %+v\n", *resp.JSON200)
}
List all Pipelines
package main
import (
"context"
"fmt"
"log"
rdpsdk "github.com/raft-tech/rdp-sdk-go"
)
func main() {
cfg, err := rdpsdk.LoadConfig()
if err != nil {
log.Fatal(err)
}
client, err := rdpsdk.NewFromConfig(cfg)
if err != nil {
log.Fatal(err)
}
resp, err := client.Pipelines().Instances().List(context.Background(), nil)
if err != nil {
log.Fatal(err)
}
if resp.JSON200 == nil {
log.Fatalf("list pipelines failed: %s", resp.Status())
}
fmt.Printf("Found %d pipelines.\n", len(*resp.JSON200))
}
Get the current status of a Pipeline
package main
import (
"context"
"encoding/json"
"fmt"
"log"
rdpsdk "github.com/raft-tech/rdp-sdk-go"
)
func main() {
cfg, err := rdpsdk.LoadConfig()
if err != nil {
log.Fatal(err)
}
client, err := rdpsdk.NewFromConfig(cfg)
if err != nil {
log.Fatal(err)
}
uid := "my-pipeline-uid"
resp, err := client.Pipelines().Instances().ID(uid).Status(context.Background())
if err != nil {
log.Fatal(err)
}
if resp.JSON200 == nil {
log.Fatalf("get pipeline status failed: %s", resp.Status())
}
status, err := json.MarshalIndent(resp.JSON200, "", " ")
if err != nil {
log.Fatal(err)
}
fmt.Printf("Status for pipeline %s:\n%s\n", uid, status)
}
Create a Pipeline
package main
import (
"context"
"fmt"
"log"
rdpsdk "github.com/raft-tech/rdp-sdk-go"
pipelinesapi "github.com/raft-tech/rdp-sdk-go/pipelines"
)
func main() {
cfg, err := rdpsdk.LoadConfig()
if err != nil {
log.Fatal(err)
}
client, err := rdpsdk.NewFromConfig(cfg)
if err != nil {
log.Fatal(err)
}
name := "example-pipeline"
description := "Example pipeline"
securityMarkings := "UNCLASSIFIED"
environment := map[string]string{"LOG_LEVEL": "info"}
transformers := []pipelinesapi.TransformerInstancePost{
{
TemplateUid: "new-transformer-abc",
Uid: "transformer",
Configuration: &pipelinesapi.Configuration{
Environment: &environment,
},
},
}
datasets := []pipelinesapi.DatasetRefPost{
{
Uid: "input-dataset",
Configuration: &pipelinesapi.DatasetConfiguration{
Ref: "quickstart-input",
ConnType: pipelinesapi.INTERNALSTREAMING,
},
},
{
Uid: "output-dataset",
Configuration: &pipelinesapi.DatasetConfiguration{
Ref: "quickstart-output",
ConnType: pipelinesapi.INTERNALOBJECTSTORE,
},
},
}
connections := []pipelinesapi.Connection{
{
From: pipelinesapi.ConnectionEndpoint{
Conn: "",
Id: "input-dataset",
Type: pipelinesapi.Dataset,
},
To: pipelinesapi.ConnectionEndpoint{
Conn: "INPUT_STREAMING_TOPIC",
Id: "transformer",
Type: pipelinesapi.Transformer,
},
},
{
From: pipelinesapi.ConnectionEndpoint{
Conn: "OUTPUT_OBJECT_BUCKET",
Id: "transformer",
Type: pipelinesapi.Transformer,
},
To: pipelinesapi.ConnectionEndpoint{
Conn: "",
Id: "output-dataset",
Type: pipelinesapi.Dataset,
},
},
}
pipeline := pipelinesapi.Pipeline{
Name: &name,
Description: &description,
SecurityMarkings: &securityMarkings,
Transformers: &transformers,
Datasets: &datasets,
Connections: &connections,
}
resp, err := client.Pipelines().Instances().Create(context.Background(), nil, pipeline)
if err != nil {
log.Fatal(err)
}
if resp.JSON201 == nil || resp.JSON201.Uid == nil {
log.Fatalf("create pipeline failed: %s", resp.Status())
}
fmt.Printf("Created Pipeline with UID %q\n", *resp.JSON201.Uid)
}
Manage a Pipeline’s state
package main
import (
"context"
"log"
rdpsdk "github.com/raft-tech/rdp-sdk-go"
)
func main() {
cfg, err := rdpsdk.LoadConfig()
if err != nil {
log.Fatal(err)
}
client, err := rdpsdk.NewFromConfig(cfg)
if err != nil {
log.Fatal(err)
}
ctx := context.Background()
pipeline := client.Pipelines().Instances().ID("my-pipeline-uid")
start, err := pipeline.Start(ctx)
if err != nil {
log.Fatal(err)
}
if start.StatusCode() >= 300 {
log.Fatalf("start pipeline failed: %s", start.Status())
}
stop, err := pipeline.Stop(ctx)
if err != nil {
log.Fatal(err)
}
if stop.StatusCode() >= 300 {
log.Fatalf("stop pipeline failed: %s", stop.Status())
}
restart, err := pipeline.Restart(ctx)
if err != nil {
log.Fatal(err)
}
if restart.StatusCode() >= 300 {
log.Fatalf("restart pipeline failed: %s", restart.Status())
}
}