Kafka (Open-Source) Event Source

For the use of Open-Source Kafka triggers, please refer to Using an Open-Source Kafka Trigger.

Kafka example event

{
  "event_version": "v1.0",
  "event_time": 1576737962,
  "trigger_type": "KAFKA",
  "region": "eu-de",
  "instance_id": "08fd3e1b-cf56-401f-b4c6-81fd2a1d3ae6",
  "records": [
      {
          "messages": [
              "kafka message1",
              "kafka message2",
              "kafka message3",
              "kafka message4",
              "kafka message5"
          ],
          "topic_id": "topic-test"
      }
  ]
}

Parameter description

Parameter

Type

Description

event_version

String

Event version

event_time

String

Time when an event occurs

trigger_type

String

Event type: KAFKA

region

String

Region where a Kafka instance resides

instance_id

String

Kafka instance ID

messages

String[]

Message content

topic_id

String

Message ID

Example

# coding: utf-8

from fg_kafkaopensource_event import KafkaOpenSourceEvent


def handler(event, context):
  logger = context.getLogger()

  logger.info("Function Name: %s", context.getFunctionName())

  kafka_opensource_event = KafkaOpenSourceEvent(event)
  
  logger.info("Trigger type: %s", kafka_opensource_event.get_trigger_type())
  
  
  output = {
    "trigger_type": kafka_opensource_event.get_trigger_type(),
  }
	
  return output

Full sample code is available in the samples-doc/scratch-event-kafkaopensource.

Package description

Public package exports for fg_kafkaopensource_event.

class fg_kafkaopensource_event.KafkaOpenSourceEvent(event)

Bases: object

Represents a Kafka Open Source event for FunctionGraph.

get_event_version()
get_event_time()
get_region()
get_trigger_type()
get_instance_id()
get_records()
to_json()
class fg_kafkaopensource_event.KafkaOpenSourceRecord(record)

Bases: object

Represents a topic record in a Kafka Open Source event.

get_topic_id()
get_messages()
to_json()
class fg_kafkaopensource_event.KafkaOpenSourceRecordMessage(record)

Bases: object

Represents a single message in a Kafka Open Source record.

get_message()
to_json()