ict-architecture

Labo 09: Transactional Outbox

Implementatie van een Message Relay in Node.js die periodiek onverzonden berichten naar RabbitMQ stuurt.

Uitvoeren

node relay.js
Start de message relay.

Bestanden & Configuraties

relay.js
javascript
const mysql = require('mysql2/promise');
const amqp = require('amqplib');
async function relayMessages() {
const connection = await mysql.createConnection({ host: 'localhost', user: 'root', database: 'microservices' });
const channel = await (await amqp.connect('amqp://localhost')).createChannel();
setInterval(async () => {
const [rows] = await connection.execute('SELECT * FROM outbox WHERE is_sent = FALSE');
for (const row of rows) {
channel.publish('domain_events', row.type, Buffer.from(JSON.stringify(row.payload)));
if (Math.random() < 0.1) process.exit(1); // Simuleer 10% falen
await connection.execute('UPDATE outbox SET is_sent = TRUE WHERE id = ?', [row.id]);
}
}, 10000);
}
relayMessages();
Periodieke controle op de outbox tabel.
outbox.sql
sql
CREATE TABLE outbox (
id INT AUTO_INCREMENT PRIMARY KEY,
aggregate_type VARCHAR(255),
aggregate_id VARCHAR(255),
type VARCHAR(255),
payload JSON,
is_sent BOOLEAN DEFAULT FALSE,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
Schema voor de outbox tabel.