ワンクリックで
rb-karafka
Building Kafka consumers and producers with Karafka in Rails: multi-threaded batch processing, message handlers, OffsetStore.
Codex または Claude でインストール この Prompt をコピーして Codex、Claude、または他のアシスタントに貼り付けると、Skill ページを確認してインストールできます。
メニュー
Building Kafka consumers and producers with Karafka in Rails: multi-threaded batch processing, message handlers, OffsetStore.
Codex または Claude でインストール この Prompt をコピーして Codex、Claude、または他のアシスタントに貼り付けると、Skill ページを確認してインストールできます。
SOC 職業分類に基づく
CONTRIBUTOR TOOL - Track CC changelog, extract new versions since last check, analyze impact on plugin (breaking changes, opportunities, deprecations). Run periodically or before releases. NOT part of the distributed plugin.
CONTRIBUTOR TOOL - Check the plugin against current cached Claude Code docs. Use before releases or after Claude docs changes to separate real schema drift from stale local assumptions. Not part of the distributed plugin.
Guide plugin development workflow for this repo. Use when editing shipped plugin files under plugins/ruby-grape-rails/, release/docs metadata, or contributor tooling under .claude/.
Analyze observational skill-effectiveness signals across scanned sessions. Use for exploratory monitoring and recommendation triage, not as a release gate.
Initializing the Ruby/Rails/Grape plugin: writes stack notes (queues, ORM-per-package, layout) into CLAUDE.md. Triggers: "initialize plugin", "setup ruby plugin", "configure Claude for Rails".
Walking a user through the Ruby/Rails/Grape plugin commands, capabilities, and workflow. Tutorial-style intro for newcomers who want to learn what the plugin offers rather than tackle a specific task.
| name | rb:karafka |
| description | Building Kafka consumers and producers with Karafka in Rails: multi-threaded batch processing, message handlers, OffsetStore. |
| effort | medium |
| disable-model-invocation | true |
Apache Kafka for Ruby with high-performance multi-threaded processing.
gem 'karafka', '~> 2.5'
gem 'karafka-web' # For Web UI
Note: Karafka 2.5+ is the current stable release. Check karafka.io for the latest version.
# karafka.rb
class KarafkaApp < Karafka::App
setup do |config|
config.kafka = {
'bootstrap.servers': ENV['KAFKA_BROKERS'],
'client.id': 'my_app'
}
config.consumer_persistence = true
end
routes.draw do
topic :orders do
consumer OrdersConsumer
end
topic :payments do
consumer PaymentsConsumer
end
end
end
KarafkaApp.boot!
class OrdersConsumer < ApplicationConsumer
def consume
# Batch processing (recommended)
Order.insert_all(messages.payloads)
mark_as_consumed(messages.last)
# Or individual processing
messages.each do |message|
process_payment(message.payload)
mark_as_consumed(message)
end
end
end
class ApplicationConsumer < Karafka::BaseConsumer
private
def mark_as_consumed(message)
consumer.mark_as_consumed(message)
end
end
# Sync (for critical messages)
Karafka.producer.produce_sync(
topic: 'orders',
payload: order.to_json,
key: order.id # Partition key
)
# Async (for high throughput)
Karafka.producer.produce_async(
topic: 'orders',
payload: order.to_json,
key: order.id
)
# Batch
messages = orders.map { |o| { topic: 'orders', payload: o.to_json, key: o.id } }
Karafka.producer.produce_many_async(messages)
class OrdersConsumer < ApplicationConsumer
def consume
messages.each do |message|
begin
process_message(message)
mark_as_consumed(message)
rescue => e
send_to_dlq(message, e)
mark_as_consumed(message)
end
end
end
private
def send_to_dlq(message, error)
Karafka.producer.produce_sync(
topic: 'orders_dlq',
payload: message.payload.merge(
error: error.message,
failed_at: Time.current
)
)
end
end
routes.draw do
topic :orders do
consumer OrdersConsumer
dead_letter_queue topic: 'orders_dlq', max_retries: 3
end
topic :orders_dlq do
consumer DlqConsumer
end
end
routes.draw do
consumer_group :orders_group do
topic :orders do
consumer OrdersConsumer
end
topic :order_updates do
consumer OrderUpdatesConsumer
end
end
consumer_group :payments_group do
topic :payments do
consumer PaymentsConsumer
end
end
end
# config/application.rb
config.active_job.queue_adapter = :karafka
Karafka::ActiveJob::JobAdapter.setup
class ProcessOrderJob < ActiveJob::Base
queue_as :orders
def perform(order_id)
Order.find(order_id).process!
end
end
# Enqueue
ProcessOrderJob.perform_later(order.id)
# config/routes.rb
mount Karafka::Web::App, at: '/karafka'
Karafka::Web.setup do |config|
config.ui.sessions.secret = ENV['KARAFKA_UI_SECRET']
end
Features:
Karafka::Admin.lag('orders_group', 'orders')
def consume
messages.each do |message|
start_time = Time.current
process_message(message)
StatsD.measure('kafka.processing_time', Time.current - start_time, tags: {
topic: topic.name,
consumer: self.class.name
})
mark_as_consumed(message)
end
end
topic :orders do
consumer OrdersConsumer
# Batch settings
max_messages 100
max_wait_time 1000
# Offset management
manual_offset_management false # Auto-commit (default)
# Error handling
dead_letter_queue topic: 'orders_dlq', max_retries: 3
end
RSpec.describe OrdersConsumer do
subject(:consumer) { described_class.new }
let(:message) do
Karafka::Messages::Message.new(
{ 'id' => 1, 'status' => 'created' },
OpenStruct.new(topic: 'orders', partition: 0, offset: 0)
)
end
it 'processes the message' do
expect { consumer.on_message(message) }
.to change(Order, :count).by(1)
end
end
FROM ruby:3.4-slim
WORKDIR /app
COPY Gemfile* ./
RUN bundle install
COPY . .
CMD ["bundle", "exec", "karafka", "server"]
apiVersion: apps/v1
kind: Deployment
metadata:
name: karafka-consumer
spec:
replicas: 3
template:
spec:
containers:
- name: consumer
image: myapp:latest
command: ["bundle", "exec", "karafka", "server"]
env:
- name: KAFKA_BROKERS
value: "kafka:9092"
references/batch-processing.md — Advanced batch processing patternsreferences/deployment-guide.md — Production deployment strategiesreferences/waterdrop-producer.md — Standalone producer configuration