【问题标题】:Kafka streams tests do not correct work close卡夫卡流测试不正确的工作关闭
【发布时间】:2018-09-25 19:14:51
【问题描述】:

我有 2 个单元测试

当我运行它们时,出现以下错误

1) 测试

   @Test
    public void simpleInsertAndOutputEventPrint() throws IOException, URISyntaxException {

        GenericRecord record = getInitialEvent();

        testDriver.pipeInput(recordFactory.create(record));
        GenericRecord result =  testDriver.readOutput(detailsEventTopic, stringDeserializer, genericAvroSerde.deserializer()).value();

        Assert.assertEquals(1,result.get("tt"));

    }

2) 测试

  @Test
  public void stateStoreSimpleInsertOutputPrint()  {
       GenericRecord record = getInitialAvayaEvent();
       testDriver.pipeInput(recordFactory.create(record));
      Packet packet1 = (Packet)  store.get("dddfdfdf");
      Assert.assertEquals("ddd",packet1.getc1()); 
  }

方法初始化

  @Before
    public void setUp() throws IOException, RestClientException, URISyntaxException {

       ...

        recordFactory = new ConsumerRecordFactory<>(initialSourceTopic,new StringSerializer(),  genericAvroSerde.serializer());
        testDriver = new TopologyTestDriver(topology, props);
        this.store = testDriver.getKeyValueStore(db);

    }

当我尝试添加下一个代码时:

  @After
    public void tearDown() {
        testDriver.close(); // Close processors after finish the tests
    }

我得到了下一个错误:

[2018-09-25 22:45:38,178] ERROR stream-thread [main] Failed to delete the state directory. (org.apache.kafka.streams.processor.internals.StateDirectory)
java.nio.file.DirectoryNotEmptyException: \tmp\kafka-streams\ks-stock-analysis-appid\0_0
    at sun.nio.fs.WindowsFileSystemProvider.implDelete(WindowsFileSystemProvider.java:266)
    at sun.nio.fs.AbstractFileSystemProvider.delete(AbstractFileSystemProvider.java:103)
    at java.nio.file.Files.delete(Files.java:1126)
    at org.apache.kafka.common.utils.Utils$2.postVisitDirectory(Utils.java:740)
    at org.apache.kafka.common.utils.Utils$2.postVisitDirectory(Utils.java:723)
    at java.nio.file.Files.walkFileTree(Files.java:2688)
    at java.nio.file.Files.walkFileTree(Files.java:2742)
    at org.apache.kafka.common.utils.Utils.delete(Utils.java:723)
    at org.apache.kafka.streams.processor.internals.StateDirectory.cleanRemovedTasks(StateDirectory.java:287)
    at org.apache.kafka.streams.processor.internals.StateDirectory.clean(StateDirectory.java:228)
    at org.apache.kafka.streams.TopologyTestDriver.close(TopologyTestDriver.java:679)
    at com.dvsts.avaya.processing.topology.TopologyKafkaStreamTest.tearDown(TopologyKafkaStreamTest.java:235)

【问题讨论】:

  • LockException 的问题应该由testDriver.close(); 中的@After 修复。您的第二个错误似乎与stackoverflow.com/questions/50602512/… 有关
  • 看来你是对的。但第二个链接并没有回答我的问题,因为我无法通过处理器 api 运行 KafkaStreams#cleanUp()
  • 您可以尝试在每次测试运行时使用不同的 application.id 作为解决方法。
  • @MatthiasJ.Sax 试过了,没解决

标签: java junit apache-kafka apache-kafka-streams


【解决方案1】:

对于测试,可以为每个创建的KTable 使用IN_MEMORY("in-memory") 存储(直接或间接,例如通过聚合);这样可以避免创建任何目录,从而不再发生错误。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-03-29
    • 1970-01-01
    • 2022-01-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-11-14
    相关资源
    最近更新 更多