Zum Inhalt

RabbitMQ als begrenzte Quelle und als Ziel

RabbitMQ ist im aktuellen Runtime-Vertrag eine begrenzte Messaging-Quelle. Ein RabbitMQ-basiertes <generate> benötigt ein explizites count. Es konsumiert höchstens so viele Nachrichten und darf weniger liefern, wenn die Queue für rabbitmq_consumer_idle_timeout_ms (Standard 250) leer bleibt. Ohne count entsteht I809, bevor ein Brokerzugriff erfolgt.

Auswahl und Parallelität

Konfiguration Worker-Policy Reihenfolge Warum
nicht gesetzt oder distribution="ordered" ein Consumer FIFO-Lieferstrom des Brokers deterministische geordnete Transformation
distribution="round_robin" angeforderte Workerzahl keine globale Ordnung höherer Quelldurchsatz durch konkurrierende Consumer

Round Robin beschreibt die RabbitMQ-Verteilung zwischen Consumern; DATAMIMIC mischt die Lieferungen nicht neu. Die Ergebnismenge darf weder Duplikate noch Lücken enthalten, die Reihenfolge über Worker hinweg ist aber nicht definiert.

Auswahlverfahren für endliche Pools passen nicht zu destruktivem Queue-Konsum. cyclic, unique, selector, iterationSelector, random, reservoir, weighted, stratified und cumulated erzeugen deshalb I947. offset>0 erzeugt I930. Diese Optionen werden nie still ignoriert.

Standardmäßig läuft ein begrenzter Streaming-Load. Erzwingt die Platform Paging, wird pageSize zur Größe jedes Consume-Fensters. count="1000" pageSize="100" führt daher zehn sequenzielle Fenster aus; Page-Offsets sind Ausführungsfortschritt und kein wahlfreier Zugriff auf die Queue.

Topologie und Acknowledgement

Die Queue-Topologie wird standardmäßig vom Broker verwaltet: allow_auto_declare_queue=false, queue_expires_ms ist nicht gesetzt und DATAMIMIC deklariert oder bindet keinen Exchange. Temporäre Queues benötigen ein explizites Opt-in sowie die beabsichtigten Werte für exclusive, auto_delete und queue_expires_ms. Das verhindert AMQP-406-Fehler durch inkompatible Neudeklaration langlebiger Queues.

Bei auto_ack=false bestätigt der Consumer eine Lieferung, nachdem die Runtime sie in die begrenzte Quellseite übernommen hat. ACK nach erfolgreichem Export-Flush ist nicht Teil dieses bounded Vertrags, sondern gehört zur Continuous-Streaming-Folgearbeit.

Publisher Confirms

Das Ziel publiziert über einen langlebigen Channel, bewahrt dessen Reihenfolge und begrenzt unbestätigte Nachrichten mit rabbitmq_publisher_max_in_flight (Standard 100). flush() wartet auf ACK, NACK und Mandatory Returns. rabbitmq_publisher_confirm_timeout_ms (Standard 60000) begrenzt Backpressure, Flush und Shutdown. Fehler laufen durch den RabbitMQ-Exporter-Fehlerkatalog statt rohe Pika-Fehler offenzulegen.

Das Verbindungsprofil besitzt Host, Virtual Host und Zugangsdaten. Deployment oder Test-Harness injizieren eine eindeutige Property rabbitmq_queue und legen diese brokerverwaltete Queue vor der Ausführung an.

rabbitmq-roundtrip.xml
 1
 2
 3
 4
 5
 6
 7
 8
 9
10
11
12
13
14
<setup numProcess="1">
    <rabbitmq-exporter id="orders_out" system="rabbit_profile"
                       exchange="" routing_key="{rabbitmq_queue}"/>
    <rabbitmq-importer id="orders_in" system="rabbit_profile"
                       queue="{rabbitmq_queue}" auto_ack="false" prefetch_count="100"/>

    <generate name="created_orders" count="100" target="orders_out">
        <key name="id" generator="IncrementGenerator"/>
        <key name="status" constant="created"/>
    </generate>

    <generate name="received_orders" source="orders_in" count="100"
              distribution="ordered" target="LogExporter"/>
</setup>

Status von Continuous Consume

Auch Kafka stellt aktuell nur consume_iter_message(count) bereit. ADR-034 beschreibt einen künftigen Continuous-Pfad, ist aber weiterhin vorgeschlagen; consume_stream, Commit nach Export-Flush, Rebalance, Backpressure und geordnetes Beenden sind nicht implementiert. Kafka Continuous Consume muss zuerst entstehen. Ein gemeinsames Streaming-SPI sollte erst aus beiden konkreten Broker-Verträgen extrahiert werden.