Token导航 LogoToken导航TokenDH.com
待分类external-servicegithub未标认证来源可访问许可证需确认审计通过

spring-kafkaspring Kafka 测试

Agent Skill

用于辅助 Java 项目开发、面向对象设计、Spring 生态、Maven 或 Gradle 依赖和后端工程实践。它适合让 Agent 分析类结构、设计接口、整理服务分层、生成测试或检查常见代码坏味道。使用时需要结合项目已有架构、包结构和依赖版本,不应只按通用教程改代码;涉及数据库、事务、并发或框架配置时,应先确认运行环境和回归测试范围。

总安装

703

周安装

29

GitHub Stars

12

下载量

230
CodexClaudeCursorGemini CLI

安装说明

本站只整理中文说明和来源信息,不托管安装包,也不代用户安装。

GitHub

来源数

2

许可证

unknown

最后核验

2026-05-01

来源状态

来源可访问

安装方式

通过对话安装

复制提示词发给支持本地命令或 Skills 的 AI 助手,先确认命令和权限,再让它执行。

请帮我安装这个 Agent Skill:spring-kafka(spring Kafka 测试)
来源仓库:https://github.com/claude-dev-suite/claude-dev-suite
仓库路径:skills/spring-kafka
安装命令:
npx skills add https://github.com/claude-dev-suite/claude-dev-suite --skill spring-kafka
安装前请先检查当前环境是否支持对应 CLI,并向我确认将要执行的命令、安装目录、联网范围和文件读写权限;确认后再执行。

命令行安装

复制命令到本机终端执行。该命令会通过 npx skills 从第三方来源获取 Skill;本站只展示命令,不托管安装包,也不自动执行。

skills.shnpx skills
npx skills add https://github.com/claude-dev-suite/claude-dev-suite --skill spring-kafka

简介

在 Spring 应用中集成 Apache Kafka 消息队列,实现解耦通信。

  • 适用于日志聚合、事件溯源或微服务间异步通知等场景。
  • 支持消费者组管理、偏移量控制和消息幂等处理机制。
  • 需关注分区策略、副本同步和生产者 ACK 级别配置。
  • 技能分类暂标记为待分类,具体用途待进一步明确。spring-kafka 属于待分类类 Skill,可作为该场景下的辅助能力补充。

SKILL.md

Spring Kafka - Quick Reference

Deep Knowledge: Use mcp__documentation__fetch_docs with technology: kafka for comprehensive documentation.

Dependencies

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

Configuration

application.yml

spring:
  kafka:
    bootstrap-servers: localhost:9092
    consumer:
      group-id: my-group
      auto-offset-reset: earliest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      value-deserializer: org.springframework.kafka.support.serializer.JsonDeserializer
      properties:
        spring.json.trusted.packages: "com.example.dto"
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      acks: all
      properties:
        enable.idempotence: true
    listener:
      ack-mode: manual
      concurrency: 3

Producer Pattern

KafkaTemplate

@Service
@RequiredArgsConstructor
public class OrderProducer {

    private final KafkaTemplate<String, OrderEvent> kafkaTemplate;

    public void sendOrder(OrderEvent event) {
        kafkaTemplate.send("orders", event.getOrderId(), event)
            .whenComplete((result, ex) -> {
                if (ex != null) {
                    log.error("Failed to send order: {}", event.getOrderId(), ex);
                } else {
                    log.info("Order sent: {} to partition {}",
                        event.getOrderId(),
                        result.getRecordMetadata().partition());
                }
            });
    }

    // With headers
    public void sendWithHeaders(OrderEvent event, String correlationId) {
        ProducerRecord<String, OrderEvent> record = new ProducerRecord<>(
            "orders", event.getOrderId(), event);
        record.headers()
            .add("correlation-id", correlationId.getBytes())
            .add("source", "order-service".getBytes());

        kafkaTemplate.send(record);
    }
}

Transactional Producer

@Configuration
public class KafkaConfig {

    @Bean
    public ProducerFactory<String, Object> producerFactory() {
        Map<String, Object> config = new HashMap<>();
        config.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        config.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "tx-");
        config.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
        return new DefaultKafkaProducerFactory<>(config);
    }

    @Bean
    public KafkaTemplate<String, Object> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }

    @Bean
    public KafkaTransactionManager<String, Object> kafkaTransactionManager() {
        return new KafkaTransactionManager<>(producerFactory());
    }
}

@Service
@Transactional("kafkaTransactionManager")
public class TransactionalProducer {

    public void sendMultiple(List<OrderEvent> events) {
        events.forEach(e -> kafkaTemplate.send("orders", e.getId(), e));
    }
}

Consumer Patterns

Basic @KafkaListener

@Service
@RequiredArgsConstructor
public class OrderConsumer {

    @KafkaListener(topics = "orders", groupId = "order-processor")
    public void consume(
            @Payload OrderEvent event,
            @Header(KafkaHeaders.RECEIVED_KEY) String key,
            @Header(KafkaHeaders.RECEIVED_PARTITION) int partition,
            @Header(KafkaHeaders.OFFSET) long offset,
            Acknowledgment ack) {

        log.info("Received order: {} from partition {} offset {}",
            event.getOrderId(), partition, offset);

        try {
            processOrder(event);
            ack.acknowledge();
        } catch (Exception e) {
            log.error("Failed to process order: {}", event.getOrderId(), e);
            throw e; // Will trigger retry
        }
    }
}

Batch Consumer

@KafkaListener(
    topics = "orders",
    groupId = "batch-processor",
    containerFactory = "batchKafkaListenerContainerFactory"
)
public void consumeBatch(
        List<OrderEvent> events,
        @Header(KafkaHeaders.RECEIVED_PARTITION) List<Integer> partitions,
        Acknowledgment ack) {

    log.info("Received batch of {} orders", events.size());
    processBatch(events);
    ack.acknowledge();
}

@Bean
public ConcurrentKafkaListenerContainerFactory<String, OrderEvent> batchKafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, OrderEvent> factory =
        new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setBatchListener(true);
    factory.getContainerProperties().setAckMode(AckMode.MANUAL);
    return factory;
}

Class-Level Listener

@KafkaListener(topics = "orders", groupId = "order-handler")
@Service
public class OrderHandler {

    @KafkaHandler
    public void handleCreated(OrderCreatedEvent event) {
        // Handle order created
    }

    @KafkaHandler
    public void handleUpdated(OrderUpdatedEvent event) {
        // Handle order updated
    }

    @KafkaHandler(isDefault = true)
    public void handleDefault(Object event) {
        log.warn("Unknown event type: {}", event.getClass());
    }
}

Retry Topics (Spring Kafka 3.x)

@RetryableTopic

@RetryableTopic(
    attempts = "3",
    backoff = @Backoff(delay = 1000, multiplier = 2.0, maxDelay = 10000),
    dltStrategy = DltStrategy.FAIL_ON_ERROR,
    autoCreateTopics = "true",
    topicSuffixingStrategy = TopicSuffixingStrategy.SUFFIX_WITH_INDEX_VALUE
)
@KafkaListener(topics = "orders", groupId = "retry-consumer")
public void consumeWithRetry(OrderEvent event, Acknowledgment ack) {
    processOrder(event);
    ack.acknowledge();
}

@DltHandler
public void handleDlt(OrderEvent event,
        @Header(KafkaHeaders.RECEIVED_TOPIC) String topic,
        @Header(KafkaHeaders.EXCEPTION_MESSAGE) String errorMessage) {

    log.error("DLT received: {} from {} - error: {}",
        event.getOrderId(), topic, errorMessage);
    // Store in database for manual review
    failedOrderRepository.save(new FailedOrder(event, errorMessage));
}

Manual Retry Configuration

@Configuration
@EnableKafka
public class KafkaRetryConfig {

    @Bean
    public RetryTopicConfiguration retryTopicConfiguration(KafkaTemplate<String, Object> template) {
        return RetryTopicConfigurationBuilder
            .newInstance()
            .maxAttempts(4)
            .fixedBackOff(3000)
            .includeTopic("orders")
            .doNotAutoCreateRetryTopics()
            .create(template);
    }
}

Error Handling

Custom Error Handler

@Bean
public DefaultErrorHandler errorHandler(KafkaTemplate<String, Object> template) {
    // Send to DLT after 3 retries
    DeadLetterPublishingRecoverer recoverer = new DeadLetterPublishingRecoverer(template,
        (record, ex) -> new TopicPartition(record.topic() + ".DLT", record.partition()));

    DefaultErrorHandler handler = new DefaultErrorHandler(recoverer,
        new FixedBackOff(1000L, 3L));

    // Don't retry for these exceptions
    handler.addNotRetryableExceptions(
        ValidationException.class,
        DeserializationException.class
    );

    return handler;
}

@Bean
public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory() {
    ConcurrentKafkaListenerContainerFactory<String, Object> factory =
        new ConcurrentKafkaListenerContainerFactory<>();
    factory.setConsumerFactory(consumerFactory());
    factory.setCommonErrorHandler(errorHandler(kafkaTemplate()));
    return factory;
}

Testing

@EmbeddedKafka

@SpringBootTest
@EmbeddedKafka(partitions = 1, topics = {"orders"})
class OrderProducerTest {

    @Autowired
    private EmbeddedKafkaBroker embeddedKafka;

    @Autowired
    private OrderProducer orderProducer;

    @Test
    void shouldSendOrder() throws Exception {
        OrderEvent event = new OrderEvent("123", "CREATED");

        Map<String, Object> consumerProps = KafkaTestUtils.consumerProps(
            "test-group", "true", embeddedKafka);
        ConsumerFactory<String, OrderEvent> cf = new DefaultKafkaConsumerFactory<>(consumerProps);
        Consumer<String, OrderEvent> consumer = cf.createConsumer();
        embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "orders");

        orderProducer.sendOrder(event);

        ConsumerRecord<String, OrderEvent> record = KafkaTestUtils.getSingleRecord(consumer, "orders");
        assertThat(record.value().getOrderId()).isEqualTo("123");
    }
}

Consumer Test with CountDownLatch

@SpringBootTest
@EmbeddedKafka(partitions = 1, topics = {"orders"})
class OrderConsumerTest {

    @SpyBean
    private OrderConsumer orderConsumer;

    @Autowired
    private EmbeddedKafkaBroker embeddedKafka;

    @Test
    void shouldConsumeAndProcessOrder() throws Exception {
        CountDownLatch latch = new CountDownLatch(1);
        doAnswer(inv -> { inv.callRealMethod(); latch.countDown(); return null; })
            .when(orderConsumer).consume(any(), any());

        Map<String, Object> props = KafkaTestUtils.producerProps(embeddedKafka);
        KafkaTemplate<String, String> template = new KafkaTemplate<>(
            new DefaultKafkaProducerFactory<>(props));
        template.send("orders", "{\"orderId\":\"456\",\"status\":\"CREATED\"}");

        assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
        verify(orderConsumer).consume(any(), any());
    }
}

Testcontainers

@SpringBootTest
@Testcontainers
class OrderIntegrationTest {

    @Container
    @ServiceConnection
    static KafkaContainer kafka = new KafkaContainer(
        DockerImageName.parse("apache/kafka-native:3.8.0"));

    @Autowired
    private KafkaTemplate<String, OrderEvent> kafkaTemplate;

    @Autowired
    private OrderRepository orderRepository;

    @Test
    void shouldProcessOrderEndToEnd() throws Exception {
        kafkaTemplate.send("orders", "key-1",
            new OrderEvent("789", "CREATED")).get(10, TimeUnit.SECONDS);

        await().atMost(Duration.ofSeconds(10))
            .untilAsserted(() -> {
                Optional<Order> order = orderRepository.findById("789");
                assertThat(order).isPresent();
                assertThat(order.get().getStatus()).isEqualTo("CREATED");
            });
    }
}
Deep dive: For MockConsumer/MockProducer, @EmbeddedKafka advanced patterns, and Node.js/Python Kafka testing, see the messaging-testing-kafka skill.

Best Practices

DoDon't
Use acks=all for durabilityUse acks=0 in production
Enable idempotenceIgnore duplicate messages
Configure DLT for failuresSilently drop failed messages
Use manual acknowledgmentAuto-commit without processing
Set proper deserializer trustTrust all packages

Production Checklist

  • acks=all configured
  • Idempotence enabled
  • Consumer group ID set
  • Manual acknowledgment mode
  • Retry topics configured
  • DLT handler implemented
  • Error handler configured
  • Proper serializers set
  • Trusted packages configured
  • Monitoring metrics exposed

When NOT to Use This Skill

  • Raw Kafka - Use kafka skill for broker config
  • RabbitMQ - Use spring-amqp instead
  • Simple messaging - Consider Spring Events
  • Kafka Streams - May need additional skill

Anti-Patterns

Anti-PatternProblemSolution
Auto commitMessage lossUse manual ack
No error handlerSilent failuresConfigure error handler
No DLTLost failed messagesAdd dead letter topic
Blocking in listenerConsumer lagUse async processing
Wrong deserializerErrors on consumeMatch producer serializer
No idempotencyDuplicate processingImplement idempotent consumer

Quick Troubleshooting

ProblemDiagnosticFix
Consumer not receivingCheck group.idVerify consumer group
Serialization errorCheck value typeConfigure correct deserializer
Rebalancing oftenCheck session.timeoutIncrease timeout
Consumer lagCheck processing timeOptimize or scale consumers
Messages in DLTCheck error logsFix processing error

Reference Documentation

适合场景

01

用户想查找某类 Agent Skill 时

02

需要根据任务场景推荐可安装能力包时

03

需要对比不同来源的安装命令和来源信息时

能力概览

能力 1

按任务关键词查找相关 Skills

能力 2

展示可复制的安装命令

能力 3

保留来源站点、仓库和原始说明,方便继续核验

能力 4

展示第三方安全扫描或审计结果

安装后应在对应宿主中按原始 README 的触发条件使用;具体调用方式请以来源页面和 README 为准。

平台分布

Codex

34.84%
按下载量换算80

Claude

31.7%
按下载量换算73

Cursor

16.99%
按下载量换算39

Gemini CLI

9.89%
按下载量换算23

安全审计

Gen Agent Trust Hub

通过

Socket

通过

Snyk

通过

权限和风险

external-service

该 Skill 可能调用第三方服务、云服务或外部模型 API,使用前需要确认账号、额度、数据发送范围和服务条款。

安装前确认

本站仅展示第三方公开信息,不托管安装包,不提供自动安装或运行环境。安装前应自行审查源码、依赖和命令行为。当前只有一个来源,正式发布前建议补源仓库或其他目录站核验。

来源信息

继续浏览同类 Skills