Skip to content

Repository files navigation

Elixir Client for Apache Pulsar

CI Coverage Status Package Version hexdocs.pm

An Elixir client for Apache Pulsar.

Features

  • ⭐ 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)

Installation

Add :pulsar_elixir to your dependencies in mix.exs:

def deps do
  [
    {:pulsar, "~> 3.0.1", hex: :pulsar_elixir} <!-- x-release-please-version -->
  ]
end

Upgrading 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.

Quick Start

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
end

Put 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>"
  ]
]

Testing

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

About

An Elixir client for Apache Pulsar.

Topics

Resources

Stars

4 stars

Watchers

2 watching

Forks

Releases

Packages

Used by

Contributors

Languages