An Elixir client for Apache Pulsar.
- ⭐ Works with Broadway (see off_broadway_pulsar)
- 🔐 SSL encryption
- 🔑 Authentication
- 📥 Consumers (including stream-friendly Reader interface; see guide)
- 📤 Producers
- Partitioned topics
- 🏎️ Compacted topics
- ✅ ACK, NACK, and redelivery (see guide)
- 📝 Schemas (see guide)
- 🔪 Chunking (see guide)
- 🍱 Batching (see guide)
- 🗜️ Compression
- ☠️ Dead-letter topics (see guide)
Add :pulsar_elixir to your dependencies in mix.exs:
def deps do
[
{:pulsar, "~> 3.0.1", hex: :pulsar_elixir} <!-- x-release-please-version -->
]
endUpgrading from 2.x? See the
upgrade guide. 3.0 replaces
config :pulsar with a client in your supervision tree, moves the API onto Pulsar.Client,
Pulsar.Consumer and Pulsar.Producer, downcases the option atoms, and changes how partition
keys are routed and chunks are framed.
Assuming you have Pulsar running on localhost:6650, the quickest way to consume messages
from a Pulsar topic is using the Reader interface, reading through a client:
{:ok, _pid} = Pulsar.Client.start_link(host: "pulsar://localhost:6650")
"persistent://my-tenant/my-namespace/my-topic"
|> Pulsar.Reader.stream(timeout: 100)
|> Enum.map(fn msg -> String.to_integer(msg.payload) end)
|> Enum.filter(fn n -> rem(n, 2) == 0 end)
|> Enum.map(fn n -> n * 2 end)For more complex scenarios and assuming that you have implemented a basic consumer like the one below:
defmodule MyPulsarConsumer do
use Pulsar.Consumer.Callback
def handle_message(message, state) do
IO.puts("Received: #{message.payload}")
{:ok, state}
end
endPut a client in your supervision tree and declare its consumers and producers on it:
children = [
{Pulsar.Client,
host: "pulsar://localhost:6650",
producers: [
[topic: "persistent://my-tenant/my-namespace/my-topic", name: :my_producer]
],
consumers: [
[topic: "persistent://my-tenant/my-namespace/my-topic",
subscription_name: "my-subscription",
callback_module: MyPulsarConsumer]
]}
]
Supervisor.start_link(children, strategy: :one_for_one)The client is the only thing your tree holds; consumers and producers run under it. Sets
only known at runtime are added with Pulsar.Consumer.start/1 and Pulsar.Producer.start/1.
Resource initialization is asynchronous, so operations may temporarily return
{:error, :not_ready}. Call Pulsar.Consumer.await_ready/2 or
Pulsar.Producer.await_ready/2 when work must wait for topic discovery and worker
initialization:
:ok = Pulsar.Producer.await_ready(:my_producer, timeout: 10_000)Sending a message using the configured producer can be done as follows:
Pulsar.Producer.send(:my_producer, "Hello, Pulsar!")In a script or an IEx session, start the client directly and add to it as you go:
{:ok, _pid} = Pulsar.Client.start_link(host: "pulsar://localhost:6650")
{:ok, _pid} = Pulsar.Producer.start(topic: "persistent://public/default/t", name: :p)Brokers, consumers and producers belong to the :default client unless told otherwise.
Several clients can coexist, which is useful when connecting to more than one cluster:
children = [
{Pulsar.Client, name: :client_1, host: "pulsar://host.cluster1.com:6650"},
{Pulsar.Client, name: :client_2, host: "pulsar://host.cluster2.com:6650"}
]
Supervisor.start_link(children, strategy: :one_for_one)A consumer or producer added at runtime selects its client with :client:
Pulsar.Producer.start(
client: :client_1,
topic: "persistent://my-tenant/my-namespace/my-topic",
name: :my_producer_1
)See the architecture guide for ownership, resource lifecycle, and recovery details.
If your Pulsar cluster requires authentication, you can configure it in the client
using the auth key:
auth: [
type: Pulsar.Auth.OAuth2,
opts: [
client_id: "<YOUR-OAUTH2-CLIENT-ID>",
client_secret: "<YOUR-OAUTH2-CLIENT-SECRET>",
site: "<YOUR-OAUTH2-ISSUER-URL>",
audience: "<YOUR-OAUTH2-AUDIENCE>"
]
]Important
Do not forget to add the following line to your /etc/hosts file before running the tests:
127.0.0.1 broker1 broker2
To run the tests, run the following command:
mix test
If you want to run only a subset of tests, specify the file including the tests you want to run
mix test test/integration/consumer_test.exs
You can also run individual tests by passing the line number where they are defined
mix test test/integration/consumer_test.exs:43
The examples directory includes a number of examples that demonstrate the use of the Pulsar client.
For example:
mix run examples/bingo.exs