1"use strict";(self.webpackChunkconduit_site=self.webpackChunkconduit_site||[]).push([[1671],{28453:(e,n,o)=>{o.d(n,{R:()=>i,x:()=>c});var r=o(96540);const s={},t=r.createContext(s);function i(e){const n=r.useContext(t);return r.useMemo(function(){return"function"==typeof e?e(n):{...n,...e}},[n,e])}function c(e){let n;return n=e.disableParentContext?"function"==typeof e.components?e.components(s):e.components||s:i(e.components),r.createElement(t.Provider,{value:n},e.children)}},35743:(e,n,o)=>{o.r(n),o.d(n,{assets:()=>d,contentTitle:()=>c,default:()=>h,frontMatter:()=>i,metadata:()=>r,toc:()=>a});const r=JSON.parse('{"id":"developing/processors/building","title":"Build your own","description":"You can build your own Conduit standalone processors in Go using the Processor SDK.","source":"@site/docs/2-developing/1-processors/1-building.mdx","sourceDirName":"2-developing/1-processors","slug":"/developing/processors/building","permalink":"/docs/developing/processors/building","draft":false,"unlisted":false,"editUrl":"https://github.com/conduitio/conduitio.github.io/edit/main/docs/2-developing/1-processors/1-building.mdx","tags":[],"version":"current","lastUpdatedAt":1751447712000,"sidebarPosition":1,"frontMatter":{"title":"Build your own"},"sidebar":"docsSidebar","previous":{"title":"Conduit Processor Template","permalink":"/docs/developing/processors/conduit-processor-template"},"next":{"title":"How it works","permalink":"/docs/developing/processors/how-it-works"}}');var s=o(74848),t=o(28453);const i={title:"Build your own"},c=void 0,d={},a=[{value:"Using <code>sdk.NewProcesorFunc</code>",id:"using-sdknewprocesorfunc",level:2},{value:"Using <code>sdk.Processor</code>",id:"using-sdkprocessor",level:2},{value:"Specification",id:"specification",level:3},{value:"Example without ParamGen",id:"example-without-paramgen",level:4},{value:"Example with ParamGen",id:"example-with-paramgen",level:4},{value:"Configure",id:"configure",level:3},{value:"Open",id:"open",level:3},{value:"Process",id:"process",level:3},{value:"Teardown",id:"teardown",level:3},{value:"Entrypoint",id:"entrypoint",level:3},{value:"Schemas",id:"schemas",level:2},{value:"Logging",id:"logging",level:2},{value:"Compiling the processor",id:"compiling-the-processor",level:2}];function l(e){const n={a:"a",admonition:"admonition",code:"code",em:"em",h2:"h2",h3:"h3",h4:"h4",li:"li",ol:"ol",p:"p",pre:"pre",strong:"strong",ul:"ul",...(0,t.R)(),...e.components};return(0,s.jsxs)(s.Fragment,{children:[(0,s.jsxs)(n.p,{children:["You can build your own Conduit standalone processors in Go using the ",(0,s.jsx)(n.a,{href:"https://github.com/ConduitIO/conduit-processor-sdk",children:"Processor SDK"}),"."]}),"\n",(0,s.jsxs)(n.p,{children:["We currently only provide a Go SDK for processors. However, if you'd like to use another language for writing a processor,\nfeel free to ",(0,s.jsx)(n.a,{href:"https://github.com/ConduitIO/conduit/issues/new?assignees=&labels=feature%2Ctriage&projects=&template=1-feature-request.yml&title=Processor+SDK+%3A+%3Clanguage%3E",children:"open an issue"}),"\nand request a specific language SDK. You can also read ",(0,s.jsx)(n.a,{href:"/docs/developing/processors/how-it-works",children:"how standalone processors work"}),"\nunder the hood and build an SDK yourself."]}),"\n",(0,s.jsxs)(n.p,{children:["The ",(0,s.jsx)(n.a,{href:"https://github.com/ConduitIO/conduit-processor-sdk",children:"Processor SDK"})," exposes two ways of building processors, one for\nsimple processors without configuration parameters, and another that gives you full control over the processor."]}),"\n",(0,s.jsxs)(n.h2,{id:"using-sdknewprocesorfunc",children:["Using ",(0,s.jsx)(n.code,{children:"sdk.NewProcesorFunc"})]}),"\n",(0,s.jsxs)(n.p,{children:["If the processor is very simple and can be reduced to a single function (e.g.\nno configuration needed), then we can use ",(0,s.jsx)(n.code,{children:"sdk.NewProcessorFunc()"})," to create a processor as below:"]}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:'//go:build wasm\n\npackage main\n\nimport (\n sdk "github.com/conduitio/conduit-processor-sdk"\n)\n\nfunc main() {\n sdk.Run(&sdk.NewProcessorFunc(\n sdk.Specification{Name: "simple-processor"}),\n func(ctx context.Context, rec opencdc.Record) (opencdc.Record, error) {\n // do s
1omething with the record\n return rec\n },\n )\n}\n'})}),"\n",(0,s.jsx)(n.p,{children:"However, if the processor needs configuration, or is more complicated than only one function, then we should use\nthe full processor approach."}),"\n",(0,s.jsxs)(n.h2,{id:"using-sdkprocessor",children:["Using ",(0,s.jsx)(n.code,{children:"sdk.Processor"})]}),"\n",(0,s.jsxs)(n.p,{children:["To build the full-blown processor, the SDK contains an interface called ",(0,s.jsx)(n.a,{href:"https://pkg.go.dev/github.com/conduitio/conduit-processor-sdk#Processor",children:"sdk.Processor"}),"\nthat contains some methods to be implemented. These methods are:"]}),"\n",(0,s.jsx)(n.h3,{id:"specification",children:"Specification"}),"\n",(0,s.jsxs)(n.p,{children:["Specification contains the metadata for the processor, which can be used to define how to reference the processor,\ndescribe what the processor does and the configuration parameters it expects. Here's a list of the fields in the\n",(0,s.jsx)(n.code,{children:"sdk.Specification"})," struct and their descriptions:"]}),"\n",(0,s.jsxs)(n.ul,{children:["\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"Name"}),": the name of the processor. Note that the name should be unique across all processors, as it's used to\nreference the processor in a pipeline (see ",(0,s.jsx)(n.a,{href:"/docs/using/processors/referencing",children:"referencing processors"}),")."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"Version"}),": the version of the processor. This should be a valid ",(0,s.jsx)(n.a,{href:"https://semver.org",children:"semver"})," version and needs to be\nupdated whenever the processor's behavior changes."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"Summary"}),": a short description of what the processor does (ideally a one-liner)."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"Description"}),": a more detailed description of what the processor does. This field can contain markdown."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"Author"}),": the author of the processor."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"Parameters"}),": a map of the processor's configuration parameters. Each parameter should have a name, a type, a\ndescription, and a list of validations. Note that the validations defined on a parameter are automatically executed\nin ",(0,s.jsx)(n.code,{children:"sdk.ParseConfig"})," (see ",(0,s.jsx)(n.a,{href:"#configure",children:"Configure"}),")."]}),"\n"]}),"\n",(0,s.jsxs)(n.p,{children:["Conduit also provides ",(0,s.jsx)(n.a,{href:"https://github.com/ConduitIO/conduit-commons/tree/main/paramgen",children:(0,s.jsx)(n.code,{children:"paramgen"})}),", a helpful tool that\ngenerates the ",(0,s.jsx)(n.code,{children:"Parameters"})," map from a Go struct. This allows us to create a configuration struct that contains the\nprocessor's parameters, define default values and validations using struct tags, and generate the ",(0,s.jsx)(n.code,{children:"Parameters"})," map.\nCheck out the ",(0,s.jsx)(n.a,{href:"https://github.com/ConduitIO/conduit-commons/tree/main/paramgen",children:"ParamGen readme"})," for more details."]}),"\n",(0,s.jsx)(n.admonition,{type:"info",children:(0,s.jsxs)(n.p,{children:["Note that a processor's name and version need to be unique across all processors, as they are used to\n",(0,s.jsx)(n.a,{href:"/docs/using/processors/referencing",children:"reference"})," the processor in a pipeline. If two processors have the same name and version,\nConduit will refuse to load them."]})}),"\n",(0,s.jsx)(n.h4,{id:"example-without-paramgen",children:"Example without ParamGen"}),"\n",(0,s.jsxs)(n.p,{children:["You can define the ",(0,s.jsx)(n.code,{children:"Specification"})," method as below and manually define the parameters map:"]}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:'package example\n\nimport (\n\t"context"\n\n\t"github.com/conduitio/conduit-commons/config"\n\tsdk "github.com/conduitio/conduit-processor-sdk"\n)\n\nfunc (p *AddFieldProcessor) Specification(context.Context) (sdk.Specification, error) {\n\treturn sdk.Specification{\n\t\tName: "myAddFieldProcessor",\n\t\tSummary: "Add a field to the record.",\n\t\tDescription: `This processor lets you configure a field that will\nbe added to the record
1into field. If the payload is not\nstructured data, this processor will panic.`,\n\t\tVersion: "v1.0.0",\n\t\tAuthor: "John Doe",\n\t\tParameters: map[string]config.Parameter{\n\t\t\t"field": {\n\t\t\t\tType: config.ParameterTypeString,\n\t\t\t\tDescription: "Field is the target field that will be set.",\n\t\t\t\tValidations: []config.Validation{\n\t\t\t\t\tconfig.ValidationRequired{},\n\t\t\t\t},\n\t\t\t},\n\t\t\t"name": {\n\t\t\t\tType: config.ParameterTypeString,\n\t\t\t\tDescription: "Name is the value of the field to add.",\n\t\t\t\tValidations: []config.Validation{\n\t\t\t\t\tconfig.ValidationRequired{},\n\t\t\t\t},\n\t\t\t},\n\t\t},\n\t}, nil\n}\n'})}),"\n",(0,s.jsx)(n.h4,{id:"example-with-paramgen",children:"Example with ParamGen"}),"\n",(0,s.jsxs)(n.ol,{children:["\n",(0,s.jsx)(n.li,{children:"Add a struct that contains the needed parameters:"}),"\n"]}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:'//go:generate paramgen -output=addField_paramgen.go addFieldConfig\n\ntype addFieldConfig struct {\n\t// Field is the target field that will be set.\n\tField string `json:"field" validate:"required"`\n\t// Name is the value of the field to add.\n\tName string `json:"value" validate:"required"`\n}\n'})}),"\n",(0,s.jsxs)(n.ol,{start:"2",children:["\n",(0,s.jsx)(n.li,{children:"Generate the parameters by running:"}),"\n"]}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{children:"paramgen -output=addField_paramgen.go addFieldConfig\n"})}),"\n",(0,s.jsxs)(n.p,{children:["This will generate a file\ncalled ",(0,s.jsx)(n.code,{children:"addField_paramgen.go"})," that contains the generated parameters map, which in turn can be used under ",(0,s.jsx)(n.code,{children:"specification"}),"\nto make it simpler and shorter, example:"]}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:'//go:generate paramgen -output=addField_paramgen.go addFieldConfig\n\ntype addFieldConfig struct {\n\t// Field is the target field that will be set.\n\tField string `json:"field" validate:"required"`\n\t// Name is the value of the field to add.\n\tName string `json:"value" validate:"required"`\n}\n\nfunc (p *AddFieldProcessor) Specification(context.Context) (sdk.Specification, error) {\n\treturn sdk.Specification{\n\t\tName: "myAddFieldProcessor",\n\t\tSummary: "Add a field to the record.",\n\t\tDescription: `This processor lets you configure a field that will\nbe added to the record into field. If the payload is not\nstructured data, this processor will panic.`,\n\t\tVersion: "v1.0.0",\n\t\tAuthor: "John Doe",\n\t\tParameters: addFieldConfig{}.Parameters(), // generated by paramgen\n\t}, nil\n}\n'})}),"\n",(0,s.jsx)(n.h3,{id:"configure",children:"Configure"}),"\n",(0,s.jsx)(n.p,{children:"Configure is the first function to be called in a processor. It provides the processor with the configuration that needs\nto be validated and stored to be used in other methods.\nThis method should not open connections or any other resources. It should solely focus on parsing and validating the\nconfiguration itself."}),"\n",(0,s.jsxs)(n.p,{children:["To add custom validations, simply validate the parameters manually under this method, and return an error if the ",(0,s.jsx)(n.code,{children:"config"}),"\nmap is not valid. On the other hand, using the utility function below would apply the builtin validations to the configuration."]}),"\n",(0,s.jsxs)(n.p,{children:["The ",(0,s.jsx)(n.a,{href:"https://github.com/ConduitIO/conduit-processor-sdk",children:"Processor SDK"})," provides some useful utility functions to help implementing this method:"]}),"\n",(0,s.jsxs)(n.ul,{children:["\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"sdk.ParseConfig"}),": used to sanitize the configuration, apply defaults, validate it using builtin validations, and copy\nthe values into the target object."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"sdk.NewReferenceResolver"}),": creates a new reference resolver from the input string. The input string is a reference\nto a field in a record, check ",(0,s.jsx)(n.a,{href:"/docs/using/processors/referencing-fields",children:"Referencing record fields"})," for more details.\nThe method will return a ",(0,s.jsx)(n.code,{children:"resolver"})," that can be used to resolve a reference to the specified field in a record and\nmanipulate that field (",(0,s.jsx)(n.code,{children:"get"}),", ",(0,s.jsx)(n.code,{children:"set"})," and ",(0,s.jsx)(n.code,{children:"delete"})," the value, or ",(0,s.jsx)(n.code,{children:"rename"})," the referenced field)."]}),"\n"]}),"\n",(0,s.jsxs)(n.p,{children:["Using these utility functions, most of the ",(0,s.jsx)(n.code,{children:"Configure"})," method implementations would look something like:"]}
1),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:'func (p *AddFieldProcessor) Configure(ctx context.Context, m map[string]string) error {\n\terr := sdk.ParseConfig(ctx, m, &p.config, addFieldConfig{}.Parameters())\n\tif err != nil {\n\t\treturn fmt.Errorf("failed to parse configuration: %w", err)\n\t}\n\n\tresolver, err := sdk.NewReferenceResolver(p.config.Field)\n\tif err != nil {\n\t\treturn fmt.Errorf("failed to parse the %q param: %w", "field", err)\n\t}\n\tp.referenceResolver = resolver\n\treturn nil\n}\n'})}),"\n",(0,s.jsx)(n.h3,{id:"open",children:"Open"}),"\n",(0,s.jsx)(n.p,{children:"This function is used to open connections, start background jobs, or initialize resources that are needed for the processor."}),"\n",(0,s.jsxs)(n.p,{children:["Note that implementing this function is ",(0,s.jsx)(n.strong,{children:(0,s.jsx)(n.em,{children:"optional"})}),"."]}),"\n",(0,s.jsx)(n.h3,{id:"process",children:"Process"}),"\n",(0,s.jsx)(n.p,{children:"Process is the main show of the processor, here we would manipulate the records received and return the processed ones."}),"\n",(0,s.jsxs)(n.p,{children:["After processing the slice of records that the function got, and if no errors occurred, it should return a slice of\n",(0,s.jsx)(n.code,{children:"sdk.ProcessedRecord"})," that matches the length of the input slice. However, if an error occurred while processing a\nspecific record, then it should be reflected in the ",(0,s.jsx)(n.code,{children:"ProcessedRecord"})," with the same index as the input record, and\nshould return the slice at that index length."]}),"\n",(0,s.jsxs)(n.p,{children:["For the interface ",(0,s.jsx)(n.code,{children:"sdk.ProcessedRecord"}),", there are three main processed records types:"]}),"\n",(0,s.jsxs)(n.ol,{children:["\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"sdk.SingleRecord"}),": is a single processed record that will continue down the pipeline."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"sdk.FilterRecord"}),": is a record that will be acked and filtered out of the pipeline."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"sdk.ErrorRecord"}),": is a record that failed to be processed and will be nacked."]}),"\n"]}),"\n",(0,s.jsx)(n.p,{children:"Example:"}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:"func (p *AddFieldProcessor) Process(ctx context.Context, records []opencdc.Record) []sdk.ProcessedRecord {\n\tout := make([]sdk.ProcessedRecord, 0, len(records))\n\tfor _, record := range records {\n\t\tresolver, err := p.referenceResolver.Resolve(&record)\n\t\tif err != nil {\n\t\t\treturn append(out, sdk.ErrorRecord{Error: err})\n\t\t}\n\t\terr = resolver.Set(p.config.Name)\n\t\tif err != nil {\n\t\t\treturn append(out, sdk.ErrorRecord{Error: err})\n\t\t}\n\t\tout = append(out, sdk.SingleRecord(record))\n\t}\n\treturn out\n}\n\n"})}),"\n",(0,s.jsxs)(n.p,{children:["Note that ",(0,s.jsx)(n.code,{children:"Process"})," should be idempotent, as it may be called multiple times with the same records (e.g. after a restart\nwhen records were not flushed)."]}),"\n",(0,s.jsx)(n.h3,{id:"teardown",children:"Teardown"}),"\n",(0,s.jsxs)(n.p,{children:["This function acts like a counterpart to ",(0,s.jsx)(n.a,{href:"#open",children:(0,s.jsx)(n.code,{children:"Open"})}),", use this function to close any open connections or resources\nthat were initialized under ",(0,s.jsx)(n.code,{children:"Open"}),"."]}),"\n",(0,s.jsxs)(n.p,{children:["Note that implementing this function is also ",(0,s.jsx)(n.strong,{children:(0,s.jsx)(n.em,{children:"optional"})}),"."]}),"\n",(0,s.jsx)(n.h3,{id:"entrypoint",children:"Entrypoint"}),"\n",(0,s.jsxs)(n.p,{children:["Since the processor will be run as a standalone ",(0,s.jsx)(n.em,{children:"WASM"})," plugin, we need to add an entrypoint to it. Also, we\nshould add a ",(0,s.jsx)(n.code,{children:"go:build"})," tag to ensure that this file is only included in the build when targeting WebAssembly."]}),"\n",(0,s.jsxs)(n.p,{children:["the entrypoint will have to be in a separate package (i.e. folder), by Go convention it's normally under\n",(0,s.jsx)(n.code,{children:"cmd/my-binary-name"}),", so it would look something like:"]}
1),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{children:".\n\u251c\u2500\u2500 my-processor.go # actual processor implementation\n\u2514\u2500\u2500 cmd\n \u2514\u2500\u2500 processor\n \u2514\u2500\u2500 main.go # entrypoint\n"})}),"\n",(0,s.jsx)(n.p,{children:"Entry point example:"}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:'//go:build wasm\n\npackage main\n\nimport (\n\tsdk "github.com/conduitio/conduit-processor-sdk"\n\t"github.com/conduitio/my-processor/example"\n)\n\nfunc main() {\n\tsdk.Run(example.NewProcessor())\n}\n'})}),"\n",(0,s.jsxs)(n.p,{children:["Check ",(0,s.jsx)(n.a,{href:"#compiling-the-processor",children:"Compiling the processor"})," for what to do next, and how to compile the processor."]}),"\n",(0,s.jsx)(n.h2,{id:"schemas",children:"Schemas"}),"\n",(0,s.jsxs)(n.p,{children:["Processors have access to the schemas in the used Schema Registry. By default, if the pipeline uses a schema\nregistry and the processor gets a record with the schema info in the ",(0,s.jsx)(n.code,{children:"Metadata"}),", then the processor will have\na middleware enabled. The middleware will decode the records before they are passed to the processor\nusing their corresponding schema from the schema registry, and encode them again after the processing is done. To\nchange this default behaviour, you can change these processor's configurations accordingly:"]}),"\n",(0,s.jsxs)(n.ul,{children:["\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"sdk.schema.decode.key.enabled"}),": Whether to decode the record key using its corresponding schema from the schema registry."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"sdk.schema.decode.payload.enabled"}),": Whether to decode the record payload using its corresponding schema from the schema registry."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"sdk.schema.encode.key.enabled"}),": Whether to encode the record key using its corresponding schema from the schema registry."]}),"\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"sdk.schema.encode.payload.enabled"}),": Whether to encode the record payload using its corresponding schema from the schema registry."]}),"\n"]}),"\n",(0,s.jsxs)(n.p,{children:[(0,s.jsx)(n.strong,{children:"Example"})," of a pipeline configuration file with these parameters:"]}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-yaml",children:"version: 2.2\npipelines:\n - id: test\n status: running\n connectors:\n - id: employees-source\n type: source\n plugin: standalone:generator\n settings:\n rate: 1\n collections.str.format.type: structured\n collections.str.format.options.id: int\n collections.str.format.options.name: string\n collections.str.format.options.admin: bool\n collections.str.operations: create\n - id: logger-dest\n type: destination\n plugin: standalone:log\n processors:\n - id: access-schema\n plugin: standalone:processor-simple\n sdk.schema.decode.key.enabled: false # disabling the default behaviour\n sdk.schema.encode.key.enabled: false # disabling the default behaviour\n"})}),"\n",(0,s.jsx)(n.p,{children:"Processors can access the Schema Registry and the schemas using two utility functions:"}),"\n",(0,s.jsxs)(n.ol,{children:["\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"schema.Create"})," : You can use this utility function to create a new schema and add it to the Schema Registry.\nThis function can be called in any of the main processor methods."]}),"\n"]}),"\n",(0,s.jsxs)(n.p,{children:[(0,s.jsx)(n.strong,{children:"Example"}),":"]}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:'func (p *exampleProcessor) Open(ctx context.Context) error {\n\t// Add a new schema to the schema registry before starting to process the records.\n\tschemaBytes := []byte(`{\n "name": "record",\n "type": "record",\n "fields": [\n {\n "name": "admin",\n "type": "boolean"\n },\n {\n "name": "id",\n "type": "int"\n },\n {\n "name": "name",\n "type": "string"\n }\n ]\n }`)\n\t_, err := schema.Create(ctx, schema.TypeAvro, "subject1", schemaBytes)\n\treturn err\n}\n'})}),"\n",(0,s.jsxs)(n.ol,{start:"2",children:["\n",(0,s.jsxs)(n.li,{children:[(0,s.jsx)(n.code,{children:"schema.Get"}),": You can use this utility function to get a schema from the Schema Registry using\nits ",(0,s.jsx)(n.code,{children:"version"})," and ",(0,s.jsx)(n.code,{children:"subject"}),". This function can be called in any of the main processor methods."]}),"\n"]}),"\n",(0,s.jsxs)(n.p,{children:[(0,s.jsx)(n.strong,{children:"Example"}),":"]}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:"func (p *exampleProcessor) Process(ctx context.Context, records []opencdc.Record) []sdk.ProcessedRecord {\n\tout := make([]sdk.ProcessedRecord, 0, len(rec
1ords))\n\tfor _, record := range records {\n\t\t// get the schema subject name from the metadata\n\t\tsubject, err := rec.Metadata.GetPayloadSchemaSubject()\n\t\tif err != nil {\n\t\t\treturn append(out, sdk.ErrorRecord{Error: err})\n\t\t}\n\t\t// get the schema version from the metadata\n\t\tversion, err := rec.Metadata.GetPayloadSchemaVersion()\n\t\tif err != nil {\n\t\t\treturn append(out, sdk.ErrorRecord{Error: err})\n\t\t}\n\t\t// get the schema using the subject and the version\n\t\tget, err := schema.Get(ctx, subject, version)\n\t\t// print the schema\n\t\tfmt.Println(string(get.Bytes))\n\t}\n\treturn out\n}\n"})}),"\n",(0,s.jsx)(n.p,{children:"Example output (the printed schema):"}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-json",children:'{\n "name": "record",\n "type": "record",\n "fields": [\n {\n "name": "admin",\n "type": "boolean"\n },\n {\n "name": "id",\n "type": "int"\n },\n {\n "name": "name",\n "type": "string"\n }\n ]\n}\n'})}),"\n",(0,s.jsx)(n.h2,{id:"logging",children:"Logging"}),"\n",(0,s.jsxs)(n.p,{children:["You can get a ",(0,s.jsx)(n.code,{children:"zerolog.logger"})," instance from the context using the ",(0,s.jsx)(n.a,{href:"https://pkg.go.dev/github.com/conduitio/conduit-processor-sdk#Logger",children:(0,s.jsx)(n.code,{children:"sdk.Logger"})}),"\nfunction. This logger is pre-configured to append logs in the format expected by Conduit."]}),"\n",(0,s.jsxs)(n.p,{children:["Keep in mind that logging in the hot path (i.e. in the ",(0,s.jsx)(n.code,{children:"Process"})," method) can have a significant impact on the performance\nof the processor, therefore we recommend using the ",(0,s.jsx)(n.code,{children:"Trace"})," level for logs that are not essential for the operation of the\nprocessor."]}),"\n",(0,s.jsx)(n.p,{children:"Example:"}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-go",children:'func (p *AddFieldProcessor) Process(ctx context.Context, records []opencdc.Record) []sdk.ProcessedRecord {\n\tlogger := sdk.Logger(ctx)\n\tlogger.Trace().Msg("Processing records")\n\t// ...\n}\n'})}),"\n",(0,s.jsx)(n.h2,{id:"compiling-the-processor",children:"Compiling the processor"}),"\n",(0,s.jsxs)(n.p,{children:["Conduit uses ",(0,s.jsx)(n.a,{href:"https://webassembly.org",children:"WebAssembly"})," to run standalone processors. This means that we need to build\nthe processor as a WebAssembly module. You can do this by setting the environment variables ",(0,s.jsx)(n.code,{children:"GOARCH=wasm"})," and ",(0,s.jsx)(n.code,{children:"GOOS=wasip1"}),"\nwhen running ",(0,s.jsx)(n.code,{children:"go build"}),". This will produce a WebAssembly module that can be used as a processor in Conduit."]}),"\n",(0,s.jsx)(n.p,{children:"So, to compile the processor, run:"}),"\n",(0,s.jsx)(n.pre,{children:(0,s.jsx)(n.code,{className:"language-sh",children:"GOARCH=wasm GOOS=wasip1 go build -o processor.wasm cmd/processor/main.go\n"})}),"\n",(0,s.jsxs)(n.p,{children:[(0,s.jsx)(n.strong,{children:(0,s.jsx)(n.em,{children:"Congratulations!"})})," Now you have a new standalone processor.\nCheck ",(0,s.jsx)(n.a,{href:"/docs/developing/processors/#where-to-put-them",children:"Standalone processors"})," for details on how to use your standalone\nprocessor in a Conduit pipeline."]}),"\n",(0,s.jsx)(n.admonition,{type:"note",children:(0,s.jsxs)(n.p,{children:["To see more standalone processor examples, check out our ",(0,s.jsx)(n.a,{href:"https://github.com/ConduitIO/conduit-processor-example",children:"example processor repository"}),"."]})})]})}function h(e={}){const{wrapper:n}={...(0,t.R)(),...e.components};return n?(0,s.jsx)(n,{...e,children:(0,s.jsx)(l,{...e})}):l(e)}}}]);
Line numbers count LF bytes from the start of the resource, as the search results do. Vendor segments are library code the classifier recognised; they are stored but not indexed. Bytes are shown as Latin1 characters, one per byte.