Docs / Connectors / Confluent

Confluent Connector

Overview

The ASAPIO Integration Add-on can connect to Confluent® or Apache Kafka® brokers, using the Confluent REST Proxy.
The ASAPIO Connector is certified by Confluent (visit confluent.io for details).

Add-on / Component nameType
ASAPIO Integration Add-on – FrameworkBase component (required)
ASAPIO Integration Add-on – Connector for Confluent®/Apache Kafka®Additional package

REST Proxy

Confluent REST Proxy for Kafka is mandatory for the connectivity and subject to a separate license. See github.com/confluentinc/kafka-rest for more details.

Key features:

  1. Certified for Confluent — visit Confluent Hub for details
  2. Supports a wide range of SAP NetWeaver-based systems, including SAP ERP, S/4HANA, BW, HCM, and many more
  3. Out-of-the-box connectivity to Confluent® Platform, Confluent® Cloud, and Kafka® (REST API v2 only)
  4. REST-based outbound communication via the Confluent Kafka REST Proxy (push)
  5. Inbound interface via the Confluent Kafka REST Proxy (pull)
  6. Choose between event-driven (single events) or batch mode (multiple events, also multiple events per REST Proxy call)
  7. Supported communication direction: Outbound, Inbound
  8. Batch mode allows multi-threading with multiple SAP work processes

Block architectures for Confluent and Kafka:

Block architecture diagram: SAP connecting to Confluent via the REST Proxy
Block architecture diagram: SAP connecting to Apache Kafka via the REST Proxy

Set-up REST-based connectivity

To establish connectivity with the Confluent® Kafka® platform, proceed with the following activities and refer to the specific sections in this documentation.

  1. Create RFC destinations to the Confluent REST Proxy in SAP system settings
  2. Set up authentication against the Confluent REST Proxy
  3. Set up a connection instance in ASAPIO Integration Add-on customizing
  4. Configure an example outbound message to test the connectivity

Create RFC destination

Create a new RFC destination of type "G" (HTTP Connection to External Server).

SM59 RFC destination Technical Settings tab — target host and port for the Confluent REST Proxy
Successful RFC connection test result

Set-up authentication to REST proxy

Prerequisites: Make sure you have either the user and password available for the REST Proxy, or that you exchanged certificates with the SAP system beforehand.

For further details on how to set up the proxy, see the Confluent documentation.

Configure authentication

While creating the RFC connection to the Confluent REST Proxy (see above), specify the authentication method.

RFC destination Logon & Security tab — Basic Authentication credentials for the REST Proxy
RFC destination Logon & Security tab — SSL certificate settings for the REST Proxy

Set-up basic settings

Activate BC-Set

Business Configuration Sets (BC-Set) contain customizing and configuration-related table entries that are not imported with the add-on.

Configure cloud adapter

Add an entry for the connector to the list of cloud adapters:

ACI Handler Class: /ASADEV/CL_ACI_KAFKA_HANDLER

Cloud adapter list showing the KAFKA entry — Handler for Apache Kafka REST Proxy

Set-up cloud codepages

Specify the codepages used in the integration:

Cloud codepages configured as part of the BC-Set

Set-up connection instance

Create the connection instance customizing that ties together the RFC destination created earlier and the cloud connector type:

Connection instance with RFC destination and Cloud Type KAFKA

Set-up Error Type Mapping

Create an entry in the Error Type Mapping section and specify at least the following mapping:

Resp. CodeMessage Type
207Success
Error Type Mapping — HTTP response code 207 mapped to Success

Set-Up Connection Values

Maintain default values for the connection to Confluent: Connections → Default values.

Default AttributeDefault Attribute Value
KAFKA_ACCEPTapplication/vnd.kafka.v2+json
KAFKA_CALL_METHODPOST
KAFKA_CONTENT_TYPEapplication/vnd.kafka.json.v2+json for JSON payloads without schema information, or application/vnd.kafka.jsonschema.v2+json for JSON payloads with schema information (configured in header attributes)
Default values on the connection instance — KAFKA_ACCEPT, KAFKA_CONTENT_TYPE, KAFKA_CALL_METHOD

Set-up Kafka protocol connectivity

The connector for Kafka and Confluent also supports connectivity without a REST Proxy, via the native Kafka protocol. The Kafka connector sends the data to an ABAP daemon, which maintains a constant connection and session to your configured Kafka broker.

Kafka broker compatibility (native protocol)

On connect, the native protocol connector queries the broker (via ApiVersions) for the version range it supports per API, and negotiates the highest version both sides support. For the connection to work, your broker must support at least the minimum version listed below for each API:

APIAPI KeyWe supportPurpose
ApiVersions18v0Initial handshake, required on every broker to negotiate the versions below
SaslHandshake17v0 – v1Negotiate the SASL mechanism (PLAIN) before authentication
SaslAuthenticate36v0 – v1SASL/PLAIN authentication (username/password)
Metadata3v0 – v8Cluster/broker/topic/partition discovery, including partition leaders
Produce0v2 – v8Sending messages (record batches) to a partition leader

These are all long-established version ranges supported by any current Confluent Platform, Confluent Cloud, or Apache Kafka broker; compatibility issues are only expected with very old (legacy) broker versions. If the broker's supported range for an API doesn't overlap with ours at all, the connection fails with an API_VERSION_MISMATCH error naming the exact API and the two ranges involved.

Create RFC Destination

Create a new RFC destination of type "G" (HTTP Connection to External Server).

RFC destination Technical Settings — target host and port for the native Kafka daemon

Add the certificates for the created destination to the certificate list selected on the Logon & Security tab:

RFC destination Logon & Security tab — SSL certificate list for the native Kafka daemon

Set-up Kafka protocol cloud adapter

Add an entry for the connector to the list of cloud adapters:

Cloud adapter list showing the S4KAFKA entry used by the native Kafka connector

Set-up connection instance

Create the connection instance customizing that ties together the RFC destination created earlier and the cloud connector type:

Connection instance for the native Kafka (S4KAFKA) connector

Default Values

Default AttributeFallback valueDescription
DAEMON_BUSY_TICK_MS50Interval in milliseconds for the next processing tick while the daemon is busy and queue entries are available. Controls how frequently messages are actively processed.
DAEMON_IDLE_SLEEP_MS300Interval in milliseconds for the next processing tick while the daemon is idle and no queue entries are available. Reduces unnecessary CPU usage during idle periods.
DAEMON_MAX_ITEM_BYTES524288Maximum allowed size of a single payload in bytes.
DAEMON_QUEUE_MAX_BYTES209715200Maximum total size of the in-memory queue in bytes. If the queue would exceed this limit, new messages are dropped due to backpressure.
DAEMON_QUEUE_MAX_ITEMS10000Maximum number of messages allowed in the local in-memory queue. Additional messages are rejected once the limit is reached.
DAEMON_QUEUE_TTL_SEC0Maximum time in seconds a message may remain in the local queue before it expires. A value of 0 means queue TTL handling is disabled.
DAEMON_SEND_TIMEOUT_MS60000Timeout in milliseconds for sent messages waiting for a producer callback.
DAEMON_TIMEOUT_CHECK_MS1000Interval in milliseconds for running timeout checks on queued and pending messages.
DAEMON_SEND_PER_TICK50Maximum number of queued messages sent per daemon tick (batch size). Values below 1 are treated as 1.
DAEMON_HANDOFF_RETRY_COUNT2Number of retry attempts when handing a batch off to the daemon fails with a retryable error (daemon not attached/temporarily unreachable).
DAEMON_HANDOFF_RETRY_DELAY_MS1000Delay in milliseconds between handoff retry attempts.
KAFKA_USERNAMEkafka_userUsername used for the Kafka connection (SASL).
Default Values for the KAFKA_PROTOCOL instance showing the daemon parameters and KAFKA_USERNAME

Set up SASL configuration

To perform authentication via SASL, you need a username and a password. These can be found in the broker settings. Note that there is usually a separate endpoint for SASL, and this must be configured.

Save the username in Default Values

For SASL authentication, a username has to be saved in Default Values.

Default Values entry KAFKA_USERNAME set to admin

Save the password in SAP Secure Store

For SASL authentication, a password has to be stored in the system's SAP Secure Store.

Enter the password in the SAP Secure Store:

Maintain the Cloud Shared Secret screen for the KAFKA_PROTOCOL cloud instance

Cluster discovery, leader-aware routing, and partitioning

Bootstrap and cluster discovery – The RFC destination configured for the native protocol connector (see Create RFC Destination above) only needs to point to a single bootstrap broker endpoint; it does not need to list every broker in the cluster. On connect, the connector sends a Metadata request to this bootstrap broker and receives the full cluster topology in return: all brokers, all topics/partitions, and the current partition leader for each partition. This metadata is cached and refreshed automatically (e.g. on connection loss or leader change).

Leader-aware routing – Once the partition leader for a given topic/partition is known, the connector opens a direct connection to that broker and sends the Produce request there, not through the bootstrap broker. If a topic's partitions are spread across multiple brokers, the connector transparently maintains one connection per broker as needed and routes each message to the correct leader. This matches the behavior of standard Kafka client libraries and requires no additional configuration beyond the single bootstrap RFC destination.

Partition assignment (KAFKA_KEY_FIELD) – When a message key is provided (via KAFKA_KEY_FIELD), the connector hashes the key (MurmurHash2, the same algorithm used by Kafka's own default partitioner) to consistently select a partition: the same key value always maps to the same partition, preserving relative ordering for that key while still allowing consumers to process different partitions in parallel. Without a key, messages are distributed across partitions at random, with no ordering guarantee.

Error Handling and Tracing

The Kafka protocol connector hands off messages to the ABAP daemon asynchronously. Successful delivery and final failures (after exhausting retries) are reported back and visible in the ACI Monitor under the Message IDs tab. Intermediate retry attempts are only visible in the daemon's own log.

ACI Monitor for the KAFKA_PROTOCOL instance, showing the Message IDs tab with per-message status

Monitor and restart the Kafka daemon

The native Kafka protocol connector runs as a persistent ABAP Daemon instance, managed through ASAPIO's own daemon administration transaction (not the generic SAP daemon monitor).

ACI Daemon Manager showing 7 configured daemon instances (6 running, 1 stopped), with Reload status, Start, Stop, Restart and Config Reset actions, and columns for Daemon Name, Instance, Status, Daemon Class, Adapter, Creator, Destination, Created At, App Server and Instance ID
ACI Daemon Manager (/ASADEV/DAEMON_CONF) – overview of all configured daemon instances and their status

The daemon automatically reconnects to the broker on connection loss; a manual restart via /ASADEV/DAEMON_CONF is normally only needed for troubleshooting.

Set-up outbound messaging

Create Message Type

Example: the examples below use the Sales Order (BUS2032) event. Choose any other suitable example if required.

For each object to be sent via ACI, you have to create a message type:

WE81 message type maintenance

Activate Message Type

The created message type has to be activated:

BD50 message type activation

Set-up additional settings in 'Header Attributes'

Configure the topic to send the events to, the fields to be used for the key, and the IDs of the key/value schemas:

Header AttributeHeader Attribute Value
KAFKA_TOPIC<topic name>, e.g. sap_demo.sales_order
KAFKA_KEY_FIELD<fields for key> (separated by ";" if multiple), e.g. VBELN;AUART
KAFKA_SCHEMA_ID<id of value schema in Schema Registry>, e.g. 4711
KAFKA_KEY_SCHEMA_ID<id of key schema in Schema Registry>, e.g. 4712

Due to limitations in the REST Proxy, you always have to specify schemas for both key and value.

Format of the key with multiple fields – When a single field is configured in KAFKA_KEY_FIELD, the raw field value is used as the Kafka message key. When multiple fields are configured (separated by ;), the resulting key is a JSON object with the SAP field names as keys, e.g.:

{
  "VBELN": "162",
  "AUART": "OR"
}

This is an ASAPIO-specific convention (there is no universal standard for composite Kafka message keys) – documented here so consumers know what to expect when parsing the key.

Confluent header attributes on outbound object — KAFKA_TOPIC, KAFKA_KEY_FIELD, KAFKA_SCHEMA_ID

Using Apache Avro

Use the format function /ASADEV/ACI_AVRO_JSON_FORMAT (together with extraction function /ASADEV/ACI_GEN_PDVIEW_EXTRACT) on the outbound object configuration to serialize messages against a schema registered in the Confluent Schema Registry.

Outbound Objects details for AVRO_JSON_MATERIAL — Extraction Func. Mod. /ASADEV/ACI_GEN_PDVIEW_EXTRACT and Format Function /ASADEV/ACI_AVRO_JSON_FORMAT
Outbound object configured with the Avro extraction and format functions

Set these header attributes on the outbound object, in addition to KAFKA_TOPIC, KAFKA_SCHEMA_ID and KAFKA_KEY_SCHEMA_ID above:

Header AttributeHeader Attribute Value
KAFKA_CONTENT_TYPEapplication/vnd.kafka.avro.v2+json
Header Attributes overview for AVRO_JSON_MATERIAL — KAFKA_CONTENT_TYPE, KAFKA_KEY_SCHEMA_ID, KAFKA_SCHEMA_ID and KAFKA_TOPIC
Header attributes on the Avro-enabled outbound object

KAFKA_REGISTRY_URL does not need to be set per outbound object – the Confluent Schema Registry URL is read from the connection level (as a default value), so it only needs to be configured once for the connection, not repeated per interface.

Default Values overview for the connection instance, showing KAFKA_REGISTRY_URL configured once at connection level
Connection-level Default Values, including KAFKA_REGISTRY_URL

The payload itself must be built in the Payload Designer (/n/ASADEV/DESIGN) before the Avro schema can be registered: each table's name and each field's PayloadFieldName must match the names used in the schema, and parent/child table relationships must mirror the nesting the schema defines. A mismatch between the two surfaces as a validation error when the outbound object runs.

Registering the schema in Kafka

Before an outbound object can use a schema ID, that schema must already exist in the Confluent Schema Registry:

  1. Open Confluent Control Center and navigate to the target topic.
  2. Go to the Schema tab.
  3. Choose Set a schema (or Edit schema if one already exists and you are evolving it).
  4. Select AVRO as the schema type.
  5. Paste in the .avsc schema definition (see structure notes below).
  6. Register/Save. Control Center returns a numeric schema ID – this is the value to put into KAFKA_SCHEMA_ID (and KAFKA_KEY_SCHEMA_ID if a separate key schema was registered).

Registering the same schema content again under a different subject reuses the same ID rather than creating a duplicate; changing the schema's content (e.g. adding a field) creates a new schema ID – outbound objects using the old ID need to be updated to the new one.

Schema structure notes

Example (Material Master, MARA + MARD):

{
  "type": "record",
  "name": "MaterialEvent",
  "fields": [
    {
      "name": "MARA",
      "type": {
        "type": "array",
        "items": {
          "type": "record",
          "name": "Mara",
          "fields": [
            { "name": "MANDT", "type": "string" },
            { "name": "MATNR", "type": "string" },
            { "name": "ERSDA", "type": ["null", {"type": "int", "logicalType": "date"}], "default": null },
            { "name": "ERNAM", "type": ["null", "string"], "default": null },
            { "name": "LVORM", "type": "boolean" },
            {
              "name": "MARD",
              "type": {
                "type": "array",
                "items": {
                  "type": "record",
                  "name": "Mard",
                  "fields": [
                    { "name": "WERKS", "type": "string" },
                    { "name": "LGORT", "type": "string" }
                  ]
                }
              }
            }
          ]
        }
      }
    }
  ]
}

Set up 'Business Object Event Linkage'

Link the configuration of the outbound object to a Business Object event:

SWE2 event linkage — Object Type BUS2032, Receiver Function Module /ASADEV/ACI_EVENTS_TRIGGER, Linkage Activated

How to use Simple Notifications

Create Outbound Object configuration

Outbound object configuration for the Simple Notify method (/ASADEV/ACI_GEN_NOTIFY_KAFKA)

Test the outbound event creation

In the example above, pick any test sales order in transaction /nVA02 and force a change event, e.g. by changing the requested delivery date on header level.

For outbound messaging, you can use and even combine the following methods:

A prerequisite for all methods is to create a message type, which is used throughout the configuration process. The following sections explain the individual options.

How to use Message Builder (Generic View Extractor)

Create Outbound Object configuration

Outbound object configuration for the Message Builder method (/ASADEV/ACI_GEN_VIEW_EXTRACTOR)

Create database view

For the data events, also configure the DB view that is used to define the extraction:

Example: Sales Order view (e.g. to be used for Sales Order (BUS2032) change events)

SE11 database view definition for the Sales Order example
Fields of the Sales Order database view

Test the outbound event creation

In the example above, pick any test sales order in transaction /nVA02 and force a change event, e.g. by changing the requested delivery date on header level.

Set-up Batch Job (Job processed messaging)

The Batch Job method allows scheduling periodic full or delta loads to Kafka without relying on SAP event linkage. This is useful for initial loads or when real-time event triggers are not feasible.

Prevent synchronous call for message type

Note

With the following settings, change pointers will be set but not sent directly.

Synchronous call for message type customizing
Unchecking Sync. On for a message type to enable Batch Job delivery

Define Variant

Creating a variant for the batch job

Schedule background job

SM36 job step using report /ASADEV/AMR_REPLICATOR with the saved variant

Test background job

ASAPIO ACI Monitor showing a completed batch job run

Set-up Packed Load (split large data)

Create Outbound Object configuration

Outbound object configuration for the Packed Load method

Create database view

Note

See also Create database view above.

For the data events, also configure the DB view that is used to define the extraction:

Example: Material master view

SE11 database view definition for the material master example
Fields of the material master database view

Set-up 'Header Attributes'

Header attributeDescriptionExample
ACI_PACK_BDCP_COMMITFlag for change pointer creation. If set, change pointers will be generated for every entry. If this flag is set, a message type has to be maintained in the outbound object. Caution: this may heavily impact performance.X
ACI_PACK_TABLEName of the table to take the key fields from. This is typically different from the DB view specified in ACI_VIEW, as we only want to build packages based on the header object, and the DB view typically contains sub-objects as well.MARA
ACI_PACK_RETRY_TIMETime in seconds. This is the duration in which the framework will attempt to get a new resource from the server group.60
ACI_PACK_WHERE_COND (optional)Condition that is applied to the table defined in ACI_PACK_TABLE.e.g. AEDAT GT '20220101'
ACI_PACK_SIZENumber of entries to send.500
ACI_PACK_KEY_LENGTHLength of the key to use from ACI_PACK_TABLE (e.g. MANDT + MATNR).13
KAFKA_KEY_FIELDName of the key field.MATERIAL_NUMBER
KAFKA_TOPICTopic name in your Confluent broker.Example.topic
Packed Load header attributes on the outbound object

Execute the initial load

Warning

Depending on the amount of data, this can stress the SAP system servers immensely. Always consult with your basis team for the correct server group to use.

Executing a Packed Load initial load with Upload Type P

Set-up Dead Letter Queues (DLQ)

With release 2510, ASAPIO introduced the option to send a message to an alternate topic (Dead Letter Queue) in two cases:

Outbound Object configuration

In the Confluent/Kafka instances, the relevant header attributes need to be configured for the outbound objects that should push messages to a Dead Letter Queue:

Header AttributeHeader Attribute Example Value
DLQ_TOPICsap_demo.dlq
DLQ_ERROR_CODES900,700,302
DLQ_RETRIES2
DLQ header attributes on the outbound object — DLQ_TOPIC, DLQ_ERROR_CODES, DLQ_RETRIES

Standard Function Modules

Function Modules: Event Linkage

Function nameDescription
/ASADEV/ACI_EVENTS_TRIGGERCustomizable trigger for event processing

Function Modules: Outbound

Mandatory Response Handler: /ASADEV/ACI_KAFKA_RESP_HANDLER (Kafka response handler method)

Extraction Func. ModuleFormatting Function
/ASADEV/ACI_GEN_NOTIFY_KAFKA (Simple notification event)—
/ASADEV/ACI_GEN_VIEW_EXTRACTOR (Dynamic data selection)/ASADEV/ACI_GEN_VIEWFRM_KAFKA (Formatting function for DB-View extraction)
/ASADEV/ACI_GEN_PDVIEW_EXTRACT (Payload Designer-based extraction)/ASADEV/ACI_AVRO_JSON_FORMAT (Avro serialization against a registered Schema Registry schema)

Function Modules: Inbound

Function nameDescription
/ASADEV/ACI_SAMPLE_IDOC_JSONInbound JSON to generic IDoc
/ASADEV/ACI_SAMPLE_IDOC_JSON2Inbound JSON to the new and improved generic IDoc
/ASADEV/ACI_JSON_TO_IDOCConverts a specially formatted JSON message to a SAP standard IDoc
/ASADEV/ACI_IDOC_JSON_AS_XMLObsolete — use /ASADEV/ACI_SAMPLE_IDOC_JSON2 instead

Set-up inbound messaging

The ASAPIO Integration Add-on is delivered with a way to pull data from the Confluent® Kafka® REST Proxy back into the SAP system. This section explains configuration and customization of the inbound function modules.

With version 9.32405 (SP09), the data pull correctly supports the host header to support load-balanced REST Proxy instances, e.g. one pull process always goes to the same REST Proxy instance.

Configure Inbound Object

Create Inbound Object configuration

Inbound object configuration — Func. name /ASADEV/ACI_SAMPLE_IDOC_JSON

Pull-based inbound offers the following header attributes, which can be maintained to receive new messages: Connections → Inbound Objects → Header Attributes.

Default AttributeExample ValueDefault
KAFKA_CONTENT_TYPE (mandatory)For the call payload. This should always be application/vnd.kafka.json.v2+json for the calls to the polling APIs.application/vnd.kafka.json.v2+json
KAFKA_DOWNLOAD_TOPIC (optional)The topic to pull data from, e.g. example.topicNo default
KAFKA_DOWNLOAD_ACCEPT (optional)Depending on the consumer format, e.g. application/vnd.kafka.avro.v2+jsonapplication/vnd.kafka.json.v2+json
KAFKA_GROUPNAME (mandatory)Name of the consumer group, e.g. exampleconsumerNo default
KAFKA_INSTANCE_NAME (mandatory)Name of the consumer instance, e.g. example_matmas_consumerNo default
KAFKA_MAX_BYTES (optional)Maximum number of bytes read in one fetch operation1000000 (1 MB)
KAFKA_TIMEOUT (optional)Maximum amount of milliseconds spent fetching records1000 (1 second)
KAFKA_CONSUMER_FORMAT (optional)The format of the consumed messages, used to convert the messages into a JSON-compatible form. The REST Proxy performs the conversion, so the SAP system always receives the message as JSON. Valid values: binary, avro, json, jsonschema, protobuf.json
KAFKA_CONSUMER_OFFSET_RESET (optional)Where the consumer group starts if it is newly created (has no impact if the consumer group already exists). Valid values: earliest, latest.earliest
KAFKA_CONSUMER_AUTO_COMMIT (optional)Whether the consumer group commits offsets automatically on fetch, or the ASAPIO add-on commits offsets after handing the data off to the configured processing FM. Valid values: true, false.false
Inbound header attributes — KAFKA_DOWNLOAD_TOPIC, KAFKA_GROUPNAME, KAFKA_INSTANCE_NAME and related values

Execute Inbound message pull

Executing the inbound message pull

Inbound processing

We recommend storing the inbound data first, without any processing logic. That way the HTTP connection can be released quickly, and processing can take place asynchronously in parallel or afterwards.

The ASAPIO-specific generic IDoc type /ASADEV/ACI_GENERIC_IDOC can be used to store the message with its payload.

Generic inbound flow: the Kafka REST Proxy delivers to the Inbound Service, which stages messages as IDoc

Note

IDocs have the advantage that they can be processed multi-threaded afterwards.

Example interface for inbound processing function modules

If you don't want to use the delivered function module that creates an IDoc, you can implement a custom function module using the following interface. The payload is received in an importing parameter IT_CONTENT:

Example function module interface for custom inbound processing, with importing parameter IT_CONTENT