Token导航 LogoToken导航TokenDH.com
研究检索只读github未标认证来源可访问许可证需确认审计提醒

messaging-testing-kafkamessaging 测试 Kafka

Agent Skill

用于辅助测试设计、自动化测试、用例整理和回归验证。它适合让 Agent 编写单元测试、端到端测试、测试计划或根据失败日志定位问题。使用时需要确认项目测试框架、运行命令和夹具数据,避免为了通过测试而改坏真实逻辑;涉及浏览器或外部服务时,应区分本地模拟、测试环境和生产环境。

总安装

533

周安装

22

GitHub Stars

12

下载量

174
CodexClaudeCursorGemini CLI

安装说明

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

GitHub

来源数

2

许可证

unknown

最后核验

2026-05-01

来源状态

来源可访问

安装方式

通过对话安装

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

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

命令行安装

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

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

简介

用于辅助测试设计、自动化测试和用例整理,适合编写单元测试或端到端测试。

  • 适用于 Codex、Claude、Cursor、Gemini CLI 中的回归验证和问题定位任务。
  • 使用时需确认项目测试框架、运行命令和夹具数据,避免为通过测试而改坏逻辑。
  • 涉及浏览器或外部服务时应区分本地模拟、测试环境与生产环境。
  • 建议结合日志分析和构建检查确保测试有效性和稳定性。

SKILL.md

Kafka Integration Testing

Quick References: See quick-ref/embedded-kafka.md for @EmbeddedKafka details, quick-ref/testcontainers-kafka.md for Testcontainers patterns.

Testing Approach Selection

ApproachSpeedFidelityBest For
@EmbeddedKafkaFast (~2s startup)High (real broker, in-process)Spring Boot unit/integration tests
MockConsumer/MockProducerInstantLow (no broker)Unit testing producer/consumer logic
Testcontainers KafkaContainerSlow (~10s startup)Highest (real Docker broker)Full integration tests, CI pipelines

Decision rule: Use @EmbeddedKafka for Spring tests by default. Use Testcontainers when you need specific Kafka versions, multi-broker clusters, or non-Spring projects. Use Mocks only for isolated unit tests.

Java/Spring: @EmbeddedKafka

Dependencies

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

Producer Test

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

    @Autowired
    private EmbeddedKafkaBroker embeddedKafka;

    @Autowired
    private OrderProducer orderProducer;

    @Test
    void shouldSendOrderEvent() {
        Map<String, Object> consumerProps = KafkaTestUtils.consumerProps(
            "test-group", "true", embeddedKafka);
        consumerProps.put(JsonDeserializer.TRUSTED_PACKAGES, "*");
        DefaultKafkaConsumerFactory<String, OrderEvent> cf =
            new DefaultKafkaConsumerFactory<>(consumerProps);
        Consumer<String, OrderEvent> consumer = cf.createConsumer();
        embeddedKafka.consumeFromAnEmbeddedTopic(consumer, "orders");

        orderProducer.sendOrder(new OrderEvent("123", "CREATED"));

        ConsumerRecord<String, OrderEvent> record =
            KafkaTestUtils.getSingleRecord(consumer, "orders", Duration.ofSeconds(10));
        assertThat(record.value().getOrderId()).isEqualTo("123");
        assertThat(record.value().getStatus()).isEqualTo("CREATED");

        consumer.close();
    }
}

Consumer Test with CountDownLatch

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

    @Autowired
    private EmbeddedKafkaBroker embeddedKafka;

    @SpyBean
    private OrderConsumer orderConsumer;

    private CountDownLatch latch = new CountDownLatch(1);

    @BeforeEach
    void setup() {
        doAnswer(invocation -> {
            invocation.callRealMethod();
            latch.countDown();
            return null;
        }).when(orderConsumer).consume(any(), any());
    }

    @Test
    void shouldConsumeOrderEvent() throws Exception {
        Map<String, Object> producerProps = KafkaTestUtils.producerProps(embeddedKafka);
        DefaultKafkaProducerFactory<String, String> pf =
            new DefaultKafkaProducerFactory<>(producerProps);
        KafkaTemplate<String, String> template = new KafkaTemplate<>(pf);

        template.send("orders", "{\"orderId\":\"456\",\"status\":\"CREATED\"}");

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

Partition Assignment Wait

// Wait for consumer to be assigned partitions before sending
ContainerTestUtils.waitForAssignment(
    listenerContainer, embeddedKafka.getPartitionsPerTopic());

Java/Spring: MockConsumer/MockProducer

MockProducer (Unit Testing)

class OrderProducerUnitTest {

    private MockProducer<String, String> mockProducer;
    private KafkaTemplate<String, String> kafkaTemplate;

    @BeforeEach
    void setup() {
        mockProducer = new MockProducer<>(true, new StringSerializer(), new StringSerializer());
        ProducerFactory<String, String> pf = new MockProducerFactory<>(mockProducer);
        kafkaTemplate = new KafkaTemplate<>(pf);
    }

    @Test
    void shouldSendToCorrectTopic() {
        kafkaTemplate.send("orders", "key", "value");

        assertThat(mockProducer.history()).hasSize(1);
        assertThat(mockProducer.history().get(0).topic()).isEqualTo("orders");
        assertThat(mockProducer.history().get(0).key()).isEqualTo("key");
    }
}

MockConsumer (Unit Testing)

class OrderConsumerUnitTest {

    private MockConsumer<String, String> mockConsumer;

    @BeforeEach
    void setup() {
        mockConsumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST);
    }

    @Test
    void shouldProcessRecords() {
        mockConsumer.assign(List.of(new TopicPartition("orders", 0)));
        mockConsumer.updateBeginningOffsets(Map.of(new TopicPartition("orders", 0), 0L));

        mockConsumer.addRecord(new ConsumerRecord<>("orders", 0, 0L, "key", "{\"orderId\":\"1\"}"));

        ConsumerRecords<String, String> records = mockConsumer.poll(Duration.ofMillis(100));
        assertThat(records.count()).isEqualTo(1);
    }
}

Java/Spring: Testcontainers KafkaContainer

With @ServiceConnection (Spring Boot 3.1+)

@SpringBootTest
@Testcontainers
class KafkaIntegrationTest {

    @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 shouldProduceAndConsumeOrder() throws Exception {
        OrderEvent event = new OrderEvent("789", "CREATED");
        kafkaTemplate.send("orders", event.getOrderId(), event).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");
            });
    }
}

With @DynamicPropertySource (pre-3.1)

@DynamicPropertySource
static void kafkaProperties(DynamicPropertyRegistry registry) {
    registry.add("spring.kafka.bootstrap-servers", kafka::getBootstrapServers);
}

Dependencies

<dependency>
    <groupId>org.testcontainers</groupId>
    <artifactId>kafka</artifactId>
    <scope>test</scope>
</dependency>

Node.js: kafkajs + Testcontainers

import { KafkaContainer } from "@testcontainers/kafka";
import { Kafka } from "kafkajs";

describe("Kafka Integration", () => {
  let container: StartedTestContainer;
  let kafka: Kafka;

  beforeAll(async () => {
    container = await new KafkaContainer("apache/kafka-native:3.8.0").start();
    kafka = new Kafka({ brokers: [container.getBootstrapServers()] });
  }, 60_000);

  afterAll(async () => {
    await container.stop();
  });

  it("should produce and consume messages", async () => {
    const admin = kafka.admin();
    await admin.connect();
    await admin.createTopics({ topics: [{ topic: "test-topic", numPartitions: 1 }] });
    await admin.disconnect();

    const producer = kafka.producer();
    await producer.connect();
    await producer.send({
      topic: "test-topic",
      messages: [{ key: "key1", value: JSON.stringify({ orderId: "123" }) }],
    });
    await producer.disconnect();

    const messages: any[] = [];
    const consumer = kafka.consumer({ groupId: "test-group" });
    await consumer.connect();
    await consumer.subscribe({ topic: "test-topic", fromBeginning: true });
    await consumer.run({
      eachMessage: async ({ message }) => {
        messages.push(JSON.parse(message.value!.toString()));
      },
    });

    await new Promise((r) => setTimeout(r, 2000));
    expect(messages).toHaveLength(1);
    expect(messages[0].orderId).toBe("123");

    await consumer.disconnect();
  });
});

Python: confluent-kafka + Testcontainers

import pytest
from testcontainers.kafka import KafkaContainer
from confluent_kafka import Producer, Consumer

@pytest.fixture(scope="module")
def kafka_container():
    with KafkaContainer("confluentinc/cp-kafka:7.6.0") as kafka:
        yield kafka

@pytest.fixture
def bootstrap_servers(kafka_container):
    return kafka_container.get_bootstrap_server()

def test_produce_and_consume(bootstrap_servers):
    producer = Producer({"bootstrap.servers": bootstrap_servers})
    producer.produce("test-topic", key="key1", value=b'{"orderId": "123"}')
    producer.flush()

    consumer = Consumer({
        "bootstrap.servers": bootstrap_servers,
        "group.id": "test-group",
        "auto.offset.reset": "earliest",
    })
    consumer.subscribe(["test-topic"])

    msg = consumer.poll(timeout=10.0)
    assert msg is not None
    assert msg.error() is None
    assert b"123" in msg.value()
    consumer.close()

Anti-Patterns

Anti-PatternProblemSolution
Not waiting for partition assignmentConsumer misses messagesUse ContainerTestUtils.waitForAssignment()
Hardcoded timeouts too shortFlaky testsUse await() with atMost() or generous timeouts
Shared topic names across testsTests interfereUse unique topic names per test or @DirtiesContext
Not closing consumers in testsResource leaks, port exhaustionAlways close in @AfterEach or try-with-resources
Using auto.offset.reset=latest in testsConsumer misses messages sent before subscriptionUse earliest for test consumers

Quick Troubleshooting

ProblemCauseSolution
"No records found"Consumer not assigned partitionsWait for assignment, use earliest offset reset
EmbeddedKafka port conflictMultiple test classes sharing brokerUse @DirtiesContext or EmbeddedKafkaHolder pattern
Deserialization errorsMismatched serializer configSet spring.json.trusted.packages in test properties
Testcontainers timeoutDocker not running or slow pullCheck Docker, increase startup timeout
Consumer lag in testsConsumer not started before produceStart consumer first, then produce

Reference Documentation

Cross-reference: For Spring Kafka producer/consumer patterns, see spring-kafka skill. For generic Testcontainers patterns, see testcontainers skill.

适合场景

01

用户想查找某类 Agent Skill 时

02

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

03

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

能力概览

能力 1

按任务关键词查找相关 Skills

能力 2

展示可复制的安装命令

能力 3

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

能力 4

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

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

平台分布

Codex

38.31%
按下载量换算67

Claude

27.32%
按下载量换算48

Cursor

18.97%
按下载量换算33

Gemini CLI

10.26%
按下载量换算18

安全审计

Gen Agent Trust Hub

通过

Socket

通过

Snyk

可疑

权限和风险

只读

该 Skill 主要提供规则、说明或参考内容,本身偏只读;真正读写文件、联网或执行命令仍取决于宿主 Agent 的任务。

安装前确认

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

来源信息

继续浏览同类 Skills