eddy147

eddy147

KafkaEx: handle_message_set is not called

When I produce messages and fetch them with

KafkaEx.fetch(topic, 0, offset: 0)

I can see the messages produced.

However, it looks like my consumer is never called:

defmodule Kafka.Consumer do
  @moduledoc false

  use KafkaEx.GenConsumer

  alias KafkaEx.Protocol.Fetch.Message

  require Logger

  def handle_message_set(message_set, state) do
    IO.inspect("Why do I never get here")
    for %Message{value: message} <- message_set do
      Logger.debug(fn -> "message: " <> inspect(message) end)
    end

    {:async_commit, state}
  end
end

it’s get started like so:

def start_child() do
    consumer_group_opts = [
      # setting for the ConsumerGroup
      heartbeat_interval: 1_000,
      # this setting will be forwarded to the GenConsumer
      commit_interval: 1_000
    ]

    consumer_group_name = Application.get_env(:kafka_ex, :consumer_group)

    child_spec = %{
      id: KafkaEx.ConsumerGroup,
      start: {
        KafkaEx.ConsumerGroup,
        :start_link,
        [Consumer, consumer_group_name, [@topic_consumer], consumer_group_opts]
      }
    }

    DynamicSupervisor.start_child(__MODULE__, child_spec)
  end

What could be the reason it never gets into handle_message_set?

Thanks! :smiley:

First 2 of 2 Posts Switch mode

eddy147

eddy147 OP

iex(3)> Kafka.ConsumerSupervisor |> DynamicSupervisor.which_children
[{:undefined, #PID<0.449.0>, :worker, [KafkaEx.ConsumerGroup]}]
eddy147

eddy147 OP

It happened to be a timing issue.
If I run the test and add :timer.sleep(10_000) I do see proof that it passed message_handler_set.

— All posts loaded —

Where Next?

Trending in Questions Top

jonnycharles
I’m in search of an Elixir library that offers PDF generation capabilities similar to Ruby’s Prawn. While there have been discussions abo...
New
spammy
I’m looking to build a personal workflow to quickly deploy web applications written in elixir/phoenix, for local consumption (ie not on t...
New
silverdr
Using Phoenix.LiveView.TagEngine as an EEx.Engine is deprecated! To compile HEEx, use Phoenix.LiveView.TagEngine.compile/2 instead. Sta...
New
dli
Before I dive in myself, did anyone successfully sprinkle Hologram into their existing LiveView app? Looking for hints regarding: Addi...
New
bottlenecked
Hi all, I wanted to ask how the community is dealing with post-release steps. Today we have Ecto migrations, which make sure that the db...
New
michallepicki
I am using Oban and occasionally, shortly after a deployment, a handful of jobs can fail because of dependency on other parts of the syst...
New
rahultumpala
Hello, I have an Elixir backend that implements a custom protocol over TCP. I want to load test the backend and assess the performance o...
New

Other Trending Topics Top

JesseHerrick
Hey, I’m Jesse and I’m the main contributor behind Dexter, a full-featured, lightning-fast Elixir LSP optimized for large codebases. It s...
New
jimsynz
Beam Bots (or just BB for short) is a framework for building fault-tolerant robotics applications in Elixir using familiar OTP patterns. ...
New
Damirados
Hello everyone. After busy few months I am happy to announce v0.1.0 of Emerge &amp; Solve. They are GUI (Emerge) and State management (S...
New
ausimian
Emily is an Elixir library that runs Nx computations on Apple’s MLX. Install it as the default Nx backend and Nx, defn, Axon, Nx.Serving,...
New
type1fool
I just stumbled on a newly redesigned elixir-lang.org. :tada: It looks like @Software_Mansion did the work, and I think it is generally a...
New
akoutmos
@hugobarauna and I (Alex Koutmos) have been hard at work on writing a book on Nerves that takes you from simply blinking LEDs to building...
New

We're in Beta

About us Mission Statement