Kafka flow with preprocessor template

json · 43 lines
1
{
2
  "summary": "Kafka flow with preprocessor template",
3
  "value": {
4
    "modules": [],
5
    "preprocessor_module": {
6
      "id": "preprocessor",
7
      "value": {
8
        "tag": "",
9
        "type": "rawscript",
10
        "content": "/**\n * General Trigger Preprocessor\n *\n * ⚠️ This function runs BEFORE the main function.\n *\n * It processes raw trigger data (e.g., MQTT, HTTP, SQS, WebSocket, Kafka) before passing it to main().\n * Common tasks:\n * - Convert binary payloads to string/JSON\n * - Extract metadata\n * - Filter messages\n * - Add timestamps/context\n *\n * The returned object determines main() parameters:\n * - {a: 1, b: 2} → main(a, b)\n * - {msg} → main(msg)\n *\n * @param event - Trigger data (e.g., MQTT, HTTP, SQS, WebSocket, Kafka)\n * @returns Processed data for main()\n */\nexport async function preprocessor(\n  event: {\n    kind: 'kafka',\n    payload: string, // base64 encoded payload\n    brokers: string[],\n    topic: string,\n    group_id: string\n  },\n) {\n  if (event.kind === 'kafka') {\n    try {\n      // Assuming the message received is a JSON value\n      const msg = atob(event.payload);\n      const data = JSON.parse(msg);\n\n      return {\n        msg,\n        data,\n      };\n    } catch (error) {\n      throw new Error(\"Failed to parse Kafka message as JSON\");\n    }\n  }\n  \n  throw new Error(`Expected kafka trigger kind, got: ${event.kind}`);\n}",
11
        "language": "bun",
12
        "input_transforms": {
13
          "event": {
14
            "type": "static"
15
          }
16
        }
17
      }
18
    }
19
  },
20
  "schema": {
21
    "type": "object",
22
    "order": [
23
      "msg",
24
      "data"
25
    ],
26
    "$schema": "https://json-schema.org/draft/2020-12/schema",
27
    "required": [
28
      "msg",
29
      "data"
30
    ],
31
    "properties": {
32
      "msg": {
33
        "type": "string",
34
        "default": "",
35
        "description": "Decoded Kafka message payload."
36
      },
37
      "data": {
38
        "type": "object",
39
        "description": "Parsed JSON data from the Kafka message."
40
      }
41
    }
42
  }
43
}
Orvanta Hub — community scripts, flows, apps & resource types.