Element <kafka-importer>¶
Zweck: Deklariert eine Kafka-Consumer-Quelle für DataMimic-Statements.
Warum: Verwende dieses Element als begrenzte Kafka-Quelle für ein generate- oder iterate-Statement.
Beispiel¶
1 | |
Entscheidungshilfe¶
Fachlicher Nutzen: Konsumiert eine begrenzte Menge von Kafka-Records als Quelle für deterministische Transformation.
-
Verwenden, wenn
- Wenn ein Modell ein Topic mit expliziter count- oder anderer unterstützter Begrenzung liest.
-
Anderen Ansatz wählen, wenn
- Wenn ein unbegrenzt laufender Consumer benötigt wird.
-
Voraussetzungen
- Liefere Broker, Topic, Consumer-Group-Policy und ein begrenztes konsumierendes Statement.
-
Alternativen
- Verwende rabbitmq-importer für begrenzten Queue-Konsum. (Siehe:
<rabbitmq-importer>)
- Verwende rabbitmq-importer für begrenzten Queue-Konsum. (Siehe:
Vollständige Beispiele¶
Eine begrenzte Kafka-Quelle mit Artefakt- und Topic-Ziel verbinden
Verwende getrennte Importer- und Exporter-IDs, damit Konsum und Veröffentlichung unabhängige Verträge mit expliziten Topics bleiben.
| kafka-roundtrip/datamimic.xml | |
|---|---|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 | |
Regeln und ungültige Kombinationen¶
I159 — Source Kafka Client Missing
Kafka source requested but client was None
Warum: Kafka source loading requires a configured Kafka client instance.
Lösung: Provide a valid Kafka client and retry.
I915 — Generate Offset Unsupported for Kafka
Warum: The offset attribute cannot be used with Kafka sources; it only applies to file-based sources.
Lösung: Remove the offset attribute, or use a file source instead of Kafka.
I164 — Source Kafka Client Unsupported API
Kafka client '{client_type}' must implement consume_iter_message(count)
Warum: The configured Kafka client does not implement the required consume API contract.
Lösung: Use a Kafka client adapter that implements consume_iter_message.
I167 — Source Kafka TLS Failed
Kafka TLS/SSL handshake failed for topic '{topic}': {reason}
Warum: The Kafka source client failed the TLS/SSL handshake while reading messages.
Lösung: Verify Kafka TLS certificate paths, trust chain, hostname validation, and protocol settings.
I168 — Source Kafka Authentication Failed
Kafka authentication failed for topic '{topic}': {reason}
Warum: The Kafka source client could not authenticate with the configured credentials.
Lösung: Verify SASL/TLS identity settings and credentials for the Kafka source client.
I169 — Source Kafka Authorization Failed
Kafka authorization failed for topic '{topic}': {reason}
Warum: The Kafka source client was authenticated but lacks permissions for the topic/resource.
Lösung: Grant required read permissions for the topic and related Kafka resources.
I170 — Source Kafka Network Failed
Kafka network connectivity failed for topic '{topic}': {reason}
Warum: The Kafka source client could not reach the configured broker endpoint.
Lösung: Verify bootstrap server host/port, DNS, routing, and firewall connectivity.
I171 — Source Kafka Timeout
Kafka read timed out for topic '{topic}': {reason}
Warum: The Kafka source operation exceeded the configured timeout while reading messages.
Lösung: Retry or tune Kafka consumer timeout settings for this workload.
I172 — Source Kafka Topic Not Found
Kafka topic/resource not found for topic '{topic}': {reason}
Warum: The configured Kafka topic/resource was not found or is not accessible.
Lösung: Verify the topic exists, is spelled correctly, and is visible to the client.
I173 — Source Kafka Read Failed
Kafka source read failed for topic '{topic}': {reason}
Warum: The Kafka source client failed while reading messages due to an unexpected runtime error.
Lösung: Review Kafka client configuration and runtime logs, then retry the operation.
I948 — Source Kafka Option Unsupported
Kafka source '{source}' does not support {option}={value}. Supported behavior: {supported}.
Warum: The requested source-selection option cannot be preserved by bounded Kafka topic consumption.
Lösung: Remove the unsupported option or use the listed Kafka source behavior.
Erlaubte Elternelemente / Erlaubte Kindelemente¶
Erlaubte Elternelemente: else, else-if, if, setup, while
Erlaubte Kindelemente:
Keine
Erweiterungsattribute: Die deklarierten Attribute sind vollständig. Dieses Element akzeptiert zusätzlich laufzeitdefinierte Herstellererweiterungen, für die es bewusst keine statische Completion gibt.
Zugehörige Konzepte und Anwendungsfälle¶
- Source und Target wählen — Guides users and agents from an input ownership boundary and intended side effect to one bounded source and target contract.
- Determinismus, Parallelität, Paging und Distribution — Separates seed, distribution, worker policy, paging, and streaming so topology choices do not become false determinism claims.
Attribute¶
Alle 56 Attribute anzeigen
allow_auto_create_topics
Allow the Kafka client to create missing topics when the broker permits it.
optional; boolean; Standardwert: true.
api_version
Kafka protocol version as underscore-separated text or an integer tuple.
optional; string; Standardwert: null.
api_version_auto_timeout_ms
Timeout in milliseconds for automatic broker-version detection.
optional; integer; Standardwert: null.
auto_commit_interval_ms
Interval in milliseconds between automatic offset commits.
optional; integer; Standardwert: null.
auto_offset_reset
Policy used when the consumer group has no valid committed offset.
optional; string; Standardwert: null; Werte: latest, earliest, none.
bootstrap_servers
Broker list (host:port).
optional; string.
check_crcs
Validate CRC checksums of consumed records.
optional; boolean; Standardwert: null.
client_id
Client identifier reported to Kafka brokers.
optional; string; Standardwert: null.
connections_max_idle_ms
Idle time in milliseconds after which a broker connection is closed.
optional; integer; Standardwert: null.
consumer_timeout_ms
Iterator timeout in milliseconds when no records are available.
optional; integer; Standardwert: null.
decoding
Optional decoding strategy.
optional; string; Standardwert: null.
enable_auto_commit
Automatically commit consumed offsets in the background.
optional; boolean; Standardwert: null.
environment
Environment label for this importer.
optional; string; Standardwert: null.
exclude_internal_topics
Exclude Kafka internal topics from subscription results.
optional; boolean; Standardwert: null.
fetch_max_bytes
Maximum bytes returned by one consumer fetch request.
optional; integer; Standardwert: null.
fetch_max_wait_ms
Maximum broker wait in milliseconds before returning a fetch response.
optional; integer; Standardwert: null.
fetch_min_bytes
Minimum bytes a broker should collect before returning a fetch response.
optional; integer; Standardwert: null.
format
Data format.
optional; string; Standardwert: null; Werte: string, json, avro.
group_id
Consumer group identifier used for offset coordination.
optional; string; Standardwert: null.
heartbeat_interval_ms
Interval in milliseconds between consumer-group heartbeats.
optional; integer; Standardwert: null.
id
Kafka importer identifier.
erforderlich; string.
max_in_flight_requests_per_connection
Maximum number of unacknowledged requests per broker connection.
optional; integer; Standardwert: null.
max_partition_fetch_bytes
Maximum bytes fetched from one partition per request.
optional; integer; Standardwert: null.
max_poll_interval_ms
Maximum interval in milliseconds between consumer polls.
optional; integer; Standardwert: null.
max_poll_records
Maximum records returned by one consumer poll.
optional; integer; Standardwert: null.
metadata_max_age_ms
Maximum age in milliseconds of cached broker metadata.
optional; integer; Standardwert: null.
metrics_num_samples
Number of samples retained for Kafka client metrics.
optional; integer; Standardwert: null.
metrics_sample_window_ms
Duration in milliseconds of each Kafka client metrics sample.
optional; integer; Standardwert: null.
page_size
Maximum number of consumed records exposed as one DataMimic source page.
optional; integer; Standardwert: null.
partition
Optional zero-based topic partition.
optional; integer; Standardwert: null.
receive_buffer_bytes
TCP receive-buffer size in bytes; use the Kafka client default when omitted.
optional; integer; Standardwert: null.
reconnect_backoff_max_ms
Maximum delay in milliseconds between broker reconnection attempts.
optional; integer; Standardwert: null.
reconnect_backoff_ms
Initial delay in milliseconds before reconnecting to a broker.
optional; integer; Standardwert: null.
request_timeout_ms
Maximum time in milliseconds to wait for a broker request.
optional; integer; Standardwert: null.
retry_backoff_ms
Delay in milliseconds before retrying a failed broker operation.
optional; integer; Standardwert: null.
sasl_kerberos_domain_name
Kerberos domain name used by GSSAPI authentication.
optional; string; Standardwert: null.
sasl_kerberos_name
Kerberos principal name used by GSSAPI authentication.
optional; string; Standardwert: null.
sasl_kerberos_service_name
Kerberos service name used by GSSAPI authentication.
optional; string; Standardwert: null.
sasl_mechanism
SASL authentication mechanism.
optional; string; Standardwert: null; Werte: PLAIN, GSSAPI, OAUTHBEARER, SCRAM-SHA-256, SCRAM-SHA-512.
sasl_plain_password
Password for PLAIN or SCRAM authentication.
optional; string; Standardwert: null.
sasl_plain_username
Username for PLAIN or SCRAM authentication.
optional; string; Standardwert: null.
schema
Schema/subject for schema-registry aware payloads.
optional; string; Standardwert: null.
security_protocol
Transport and authentication protocol used for broker connections.
optional; string; Standardwert: null; Werte: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL.
send_buffer_bytes
TCP send-buffer size in bytes; use the Kafka client default when omitted.
optional; integer; Standardwert: null.
session_timeout_ms
Consumer-group session timeout in milliseconds.
optional; integer; Standardwert: null.
socks5_proxy
SOCKS5 proxy address used for broker connections.
optional; string; Standardwert: null.
ssl_cafile
Path to the CA certificate file used to verify brokers.
optional; string; Standardwert: null.
ssl_certfile
Path to the client certificate file.
optional; string; Standardwert: null.
ssl_check_hostname
Verify that broker certificates match their host names.
optional; boolean; Standardwert: null.
ssl_cipher_suites
OpenSSL cipher-suite expression used for broker connections.
optional; string; Standardwert: null.
ssl_crlfile
Path to a certificate-revocation-list file.
optional; string; Standardwert: null.
ssl_keyfile
Path to the client private-key file.
optional; string; Standardwert: null.
ssl_password
Password used to decrypt the client private key.
optional; string; Standardwert: null.
ssl_protocol
SSL protocol name passed to the Kafka client.
optional; string; Standardwert: null.
system
System identifier for this importer.
optional; string; Standardwert: null.
topic
Topic to subscribe to.
erforderlich; string.