【问题标题】:Test stream in akka test kit using Java使用 Java 的 akka 测试套件中的测试流
【发布时间】:2019-11-21 12:02:40
【问题描述】:

https://doc.akka.io/docs/akka/current/stream/stream-testkit.html

我正在使用 Java 使用 akka,谁能告诉我系统是如何在代码中初始化的

final Sink<Integer, CompletionStage<Integer>> sinkUnderTest =
Flow.of(Integer.class)
    .map(i -> i * 2)
    .toMat(Sink.fold(0, (agg, next) -> agg + next), Keep.right());

final CompletionStage<Integer> future =
Source.from(Arrays.asList(1, 2, 3, 4)).runWith(sinkUnderTest, system);
final Integer result = future.toCompletableFuture().get(3, TimeUnit.SECONDS);
assert (result == 20);

static ActorSystem system =ActorSystem.create() 不工作

Source.from(Arrays.asList(1, 2, 3, 4)).runWith(sinkUnderTest, system);

【问题讨论】:

  • 当你说“ActorSystem.create() 不起作用”时,它是如何表现出来的?它会抛出异常还是什么?
  • 为我工作。您在运行它时遇到了什么行为以及您使用了哪些导入和依赖项?
  • 它给出了编译错误并期望一个带有 runwith 的物化器

标签: java akka akka-stream


【解决方案1】:

编辑:

当您将 akka 2.5 与 2.12 api 一起使用时,您必须按照快速入门部分中的描述创建一个物化器以获取匹配文档,请查看 here

private ActorSystem system;
private Materializer materializer;

@BeforeEach
public void setup() {
    system = ActorSystem.create("StreamTestKitDocTest");
    materializer = ActorMaterializer.create(system);
}

// ...

@Test
public void test() throws InterruptedException, ExecutionException, TimeoutException {
    // ...
    Source.from(Arrays.asList(1, 2, 3, 4)).runWith(sinkUnderTest, materializer);
    // ...
}

正如您所说,它甚至无法为您编译,我假设您的依赖项或导入存在问题。一个常见的错误是您不小心导入了 scala 版本,而不是 java dsl。

这里是我用来验证其工作的依赖项和代码(用 Java 1.8 测试):

依赖关系:

<dependencies>
    <dependency>
      <groupId>com.typesafe.akka</groupId>
      <artifactId>akka-stream-testkit_2.13</artifactId>
      <version>2.6.0</version>
      <scope>test</scope>
    </dependency>
     <dependency>
        <groupId>junit</groupId>
        <artifactId>junit</artifactId>
        <version>4.12</version>
        <scope>test</scope>
    </dependency>
</dependencies>

JUnit 测试用例:

import java.util.Arrays;
import java.util.concurrent.CompletionStage;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;

import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
import org.junit.Test;

import akka.actor.ActorSystem;
import akka.stream.javadsl.Flow;
import akka.stream.javadsl.Keep;
import akka.stream.javadsl.Sink;
import akka.stream.javadsl.Source;

public class AkkaTest {

    private ActorSystem system;

    @Before
    public void setup() {
        system = ActorSystem.create("StreamTestKitDocTest");
    }

    @After
    public void shutdown() {
        system.terminate();
    }

    @Test
    public void test() throws InterruptedException, ExecutionException, TimeoutException {

        final Sink<Integer, CompletionStage<Integer>> sinkUnderTest =
        Flow.of(Integer.class)
            .map(i -> i * 2)
            .toMat(Sink.fold(0, (agg, next) -> agg + next), Keep.right());

        final CompletionStage<Integer> future =
        Source.from(Arrays.asList(1, 2, 3, 4)).runWith(sinkUnderTest, system);
        final Integer result = future.toCompletableFuture().get(3, TimeUnit.SECONDS);

        Assert.assertEquals(20, result.intValue());
    }
}

【讨论】:

  • 我看到的是,无论是在 java dsl 还是 scala dsl 下,所有文件都在 jar 中的 *.scala 中,不知道为什么
  • 我在 gradle testImplementation 'org.junit.jupiter:junit-jupiter:5.5.2' testImplementation "com.typesafe.akka:akka-stream-testkit_2.12:2.5" implementation 中使用依赖项com.typesafe.akka:akka-stream_2.12:2.5"
  • 该版本使用不同的 API。我添加了文档的正确链接和答案的简短示例。
【解决方案2】:

你检查过示例的源代码吗? source code

  static ActorSystem system;

  @BeforeClass
  public static void setup() {
    system = ActorSystem.create("StreamTestKitDocTest");
  }

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-07-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-05-08
    • 2019-09-29
    • 2011-10-15
    相关资源
    最近更新 更多