【问题标题】:How to use Cassandra sink with TestContainers in Flink如何在 Flink 中使用 Cassandra sink 和 TestContainers
【发布时间】:2020-12-19 17:23:35
【问题描述】:

我试图在使用DataStreamTestBase 进行测试的简单Flink 管道中使用TestContainers 来测试Cassandra Sink:

public class CassandraPojoSinkExampleTest extends DataStreamTestBase {

    @Rule
    public CassandraContainer cassandra = new CassandraContainer<>().withInitScript("testInit3.cql")
            .withCreateContainerCmdModifier(cmd -> cmd.getHostConfig()
                    .withCpuCount(2L));

    @Override
    public void initialize() throws Exception {
        super.initialize();
        testEnv.getConfig().disableObjectReuse();
        testEnv.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
        testEnv.setRestartStrategy(RestartStrategies.fixedDelayRestart(1, Time.of(1, TimeUnit.SECONDS)));
    }

    @Test
    public void test() throws Throwable {
       ArrayList<Message> messages = new ArrayList<>(20);

            for (long i = 0; i < 20; i++) {
                messages.add(new Message("cassandra-" + i));
            }


        Cluster cluster = cassandra.getCluster();

        try (Session session = cluster.connect()) {

            CassandraPojoSinkExample cassandraPojoSinkExample= new CassandraPojoSinkExample();
            cassandraPojoSinkExample.doIt(cassandra.getHost(), testEnv.fromCollection(messages));

            MappingManager manager= new MappingManager(session);
            Mapper<Message> mapper=manager.mapper(Message.class);
            String word ="cassandra-1";
            Message message= mapper.get(word);
            System.out.println(message);
        }
    }

Flink 管道看起来像(它唯一的 sink 在 doIt 方法):

public class CassandraPojoSinkExample {
    private static final ArrayList<Message> messages = new ArrayList<>(20);

    static {
        for (long i = 0; i < 20; i++) {
            messages.add(new Message("cassandra-" + i));
        }
    }
    public static void main(String[] args) throws Exception {

        // the main method is not used in tests

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        DataStreamSource<Message> source = env.fromCollection(messages);

        CassandraPojoSinkExample cassandraPojoSinkExample=new CassandraPojoSinkExample();
        cassandraPojoSinkExample.doIt("127.0.0.1", source);
    }
    public void doIt(String address, DataStreamSource<Message> source) throws Exception {

        CassandraSink.addSink(source)
                .setClusterBuilder(new ClusterBuilder() {
                    @Override
                    protected Cluster buildCluster(Builder builder) {
                        return builder.addContactPoint(address).build();
                    }
                })
                .setMapperOptions(() -> new Mapper.Option[]{Mapper.Option.saveNullFields(true)})
                .build();
    }

当我运行测试时出现错误:

16:41:57,675 ERROR org.apache.flink.streaming.runtime.tasks.StreamTask           - Error during disposal of stream operator.
java.lang.NullPointerException
    at org.apache.flink.streaming.connectors.cassandra.CassandraSinkBase.checkAsyncErrors(CassandraSinkBase.java:162)
    at org.apache.flink.streaming.connectors.cassandra.CassandraSinkBase.close(CassandraSinkBase.java:96)
    at org.apache.flink.api.common.functions.util.FunctionUtils.closeFunction(FunctionUtils.java:43)
    at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.dispose(AbstractUdfStreamOperator.java:117)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.disposeAllOperators(StreamTask.java:651)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.cleanUpInvoke(StreamTask.java:562)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:480)
    at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:708)
    at org.apache.flink.runtime.taskmanager.Task.run(Task.java:533)
    at java.lang.Thread.run(Thread.java:748)
16:41:57,683 INFO  org.apache.flink.runtime.taskmanager.Task                     - Source: Collection Source -> Sink: Cassandra Sink (1/1) (721a02e3c1345d640bf03b67eff0aba7) switched from RUNNING to FAILED.
com.datastax.driver.core.exceptions.NoHostAvailableException: All host(s) tried for query failed (tried: localhost/127.0.0.1:9042 (com.datastax.driver.core.exceptions.TransportException: [localhost/127.0.0.1] Cannot connect))
    at com.datastax.driver.core.ControlConnection.reconnectInternal(ControlConnection.java:231)
    at com.datastax.driver.core.ControlConnection.connect(ControlConnection.java:77)
    at com.datastax.driver.core.Cluster$Manager.init(Cluster.java:1414)
    at com.datastax.driver.core.Cluster.init(Cluster.java:162)
    at com.datastax.driver.core.Cluster.connectAsync(Cluster.java:333)
    at com.datastax.driver.core.Cluster.connect(Cluster.java:283)
    at org.apache.flink.streaming.connectors.cassandra.CassandraPojoSink.createSession(CassandraPojoSink.java:123)
    at org.apache.flink.streaming.connectors.cassandra.CassandraSinkBase.open(CassandraSinkBase.java:87)
    at org.apache.flink.streaming.connectors.cassandra.CassandraPojoSink.open(CassandraPojoSink.java:106)
    at org.apache.flink.api.common.functions.util.FunctionUtils.openFunction(FunctionUtils.java:36)
    at org.apache.flink.streaming.api.operators.AbstractUdfStreamOperator.open(AbstractUdfStreamOperator.java:102)
    at org.apache.flink.streaming.api.operators.StreamSink.open(StreamSink.java:48)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.initializeStateAndOpen(StreamTask.java:990)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.lambda$beforeInvoke$0(StreamTask.java:453)
    at org.apache.flink.streaming.runtime.tasks.StreamTaskActionExecutor$SynchronizedStreamTaskActionExecutor.runThrowing(StreamTaskActionExecutor.java:94)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.beforeInvoke(StreamTask.java:448)
    at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:460)
    at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:708)
    at org.apache.flink.runtime.taskmanager.Task.run(Task.java:533)
    at java.lang.Thread.run(Thread.java:748)

我在这里做错了什么?

【问题讨论】:

    标签: java cassandra apache-flink testcontainers sink


    【解决方案1】:

    从上面的堆栈跟踪中,com.datastax.driver.core.exceptions.NoHostAvailableException: All host(s) tried for query failed (tried: localhost/127.0.0.1:9042 似乎 cassandra 主机不可用。

    我会说你需要将端口暴露给外部:

        @Rule
        public CassandraContainer cassandra = new CassandraContainer<>().withInitScript("testInit3.cql")
                .withCreateContainerCmdModifier(cmd -> cmd.getHostConfig()
                .withCpuCount(2L))
                .withExposedPorts(9042);
    

    另外,我会通过telnet 手动确保到相应的主机/端口,以确保服务已启动。

    【讨论】:

      【解决方案2】:

      @Mikalai Lushchytski 的回答很好,问题出在暴露的端口上。

      暴露的端口也需要在addSink方法中传递给.setHost()。

      另外,老实说,我之前已经尝试过了,但没有注意到它正在工作,因为问题在于从 Cassandra 数据库表中接收信息。在提供的代码中,我无法在测试期间从数据库中检索任何行,问题似乎与 DataStreamTestBase 和调用方法的顺序有关。我通过手动输入数据库并在正在进行的测试中提取数据来确认管道的工作。

      使用 DataStreamTestBase 从 Cassandra TestContainer 正确获取数据库记录的解决方案 是使用 Overrided executeTest() 方法:

      public class CassandraPojoSinkExampleTest extends DataStreamTestBase {
          private static boolean runTestAfterMethod = true;
      
          @Rule
          public CassandraContainer cassandra = new CassandraContainer<>()
                  .withInitScript("testInit3.cql")
                  .withExposedPorts(9042);
      
          @Before
          public void setUp() {
              runTestAfterMethod = true;
          }
      
          @Override
          public void executeTest() throws Throwable {
              if (runTestAfterMethod) {
                  super.executeTest();
              }
          }
      
          @Override
          public void initialize() throws Exception {
              super.initialize();
              testEnv.getConfig().disableObjectReuse();
              testEnv.setRestartStrategy(RestartStrategies.fixedDelayRestart(1, Time.of(1, TimeUnit.SECONDS)));
          }
      
          @Test
          public void test() throws Throwable {
              runTestAfterMethod = false;
      
              ArrayList<Message> messages = new ArrayList<>(20);
              for (long i = 0; i < 20; i++) {
                  messages.add(new Message("cassandra-" + i));
              }
      
              DataStreamSource<Message> source = testEnv.fromCollection(messages);
      
              CassandraPojoSinkExample cassandraPojoSinkExample = new CassandraPojoSinkExample();
              cassandraPojoSinkExample.doIt(cassandra.getHost(), cassandra.getMappedPort(9042), source);
      
          try( Session session = cassandra.getCluster().newSession()) {
              testEnv.executeTest();     // Crucial thing
      
              ResultSet resultSet = session.execute("select * from test.features;");
              List<Row> rows = resultSet.all();
              List<String> actualRecords = new ArrayList<String>();
      
                   if (rows.size() != 0) {
                      for (int i = 0; i < 20; i++) {
                          actualRecords.add(rows.get(0).getString(i));
                      }
                  }
      
                  List<String> expectedRecords = new ArrayList<String>() {
                      {
                          for (int i = 0; i < 20; i++) {
                              add("cassandra-" + i);
                          }
                      }
                  };
      
                  assertTrue("List equality without order",
                          actualRecords.containsAll(expectedRecords)&& actualRecords.containsAll(expectedRecords));
              }
          }
          }
      
      }
      

      这样我们就可以确定resultSet 不能为空,并且测试会以正确的顺序执行。

      【讨论】:

        猜你喜欢
        • 2019-12-23
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2019-05-01
        • 1970-01-01
        • 2022-09-23
        • 1970-01-01
        相关资源
        最近更新 更多