1. Go to this page and download the library: Download saravanasai/stream-pulse library. Choose the download type require.
2. Extract the ZIP file and open the index.php.
3. Add this code to the index.php.
<?php
require_once('vendor/autoload.php');
/* Start to develop here. Best regards https://php-download.com/ */
saravanasai / stream-pulse example snippets
// In your producer Laravel application
// Publish an event
StreamPulse::publish('orders', ['id' => 1234, 'amount' => 99.99]);
// In your consumer Laravel application
// Register an event handler - app service provider
StreamPulse::on('orders', function ($payload, $messageId) {
OrderProcessor::process($payload);
});
// Run the consumer to process events (in the consumer app)
// php artisan streampulse:consume orders
use StreamPulse\StreamPulse\Facades\StreamPulse;
// Publish an event to a topic
StreamPulse::publish('orders', [
'id' => 1234,
'customer' => 'John Doe',
'total' => 99.99,
'items' => [
['product_id' => 101, 'quantity' => 2, 'price' => 49.99]
]
]);
// Publish after DB transaction commits
StreamPulse::publishAfterCommit('orders', $orderData);
use Illuminate\Support\Facades\DB;
use StreamPulse\StreamPulse\Facades\StreamPulse;
DB::transaction(function () {
// Create an order in the database
$order = Order::create([
'customer_id' => 123,
'amount' => 99.99,
]);
// This event will only be published if the transaction commits successfully
StreamPulse::publishAfterCommit('orders', [
'id' => $order->id,
'status' => 'created',
'customer_id' => $order->customer_id,
'amount' => $order->amount,
]);
// If the transaction fails or is rolled back, no event will be published
});
// In a service provider
StreamPulse::on('orders', function ($payload, $messageId) {
// Process the order
OrderProcessor::process($payload);
});
// Then run the consumer command
// php artisan streampulse:consume orders
// Define global defaults in config/streampulse.php
'defaults' => [
'retention' => 1000, // Keep 1000 events per stream by default
],
// Override for specific topics
'topics' => [
'orders' => [
'retention' => 5000, // Keep more events for important topics
],
'logs' => [
'retention' => 500, // Keep fewer events for high-volume topics
],
],
use StreamPulse\StreamPulse\Facades\StreamPulse;
// Consume events from a topic with a consumer group
StreamPulse::consume('orders', 'order-processors', function ($payload, $messageId) {
try {
// Process the event
OrderProcessor::process($payload);
// Acknowledge the message as processed
StreamPulse::ack('orders', $messageId, 'order-processors');
} catch (\Exception $e) {
// Handle error
// The message will remain in pending state and can be retried
// After max retries, it will be moved to the DLQ automatically
}
});