Interface StreamProcessor

    • Method Summary

      All Methods Instance Methods Abstract Methods 
      Modifier and Type Method Description
      boolean drainAndStop​(Duration timeout)
      Stop computations when input streams are empty.
      Latency getLatency​(String computationName)
      Returns the latency for a computation.
      long getLowWatermark()
      Returns the low watermark for all the computations of the topology.
      long getLowWatermark​(String computationName)
      Returns the low watermark for the computation.
      StreamProcessor init​(Topology topology, Settings settings)
      Initialize streams, but don't run the computations
      boolean isDone​(long timestamp)
      Returns true if all messages with a lower timestamp has been processed by the topology.
      boolean isTerminated()
      True if there is no active processing threads.
      void shutdown()
      Shutdown immediately.
      void start()
      Run the initialized computations.
      boolean stop​(Duration timeout)
      Try to stop computations gracefully after processing a record or a timer within the timeout duration.
      String toJson​(Map<String,​String> meta)
      Returns a JSON representation of the processor, this includes the list of streams and computations with their settings along with the topology.
      boolean waitForAssignments​(Duration timeout)
      Wait for the computations to have assigned partitions ready to process records.
    • Method Detail

      • start

        void start()
        Run the initialized computations.
      • stop

        boolean stop​(Duration timeout)
        Try to stop computations gracefully after processing a record or a timer within the timeout duration. If this can not be done within the timeout, shutdown and returns false.
      • drainAndStop

        boolean drainAndStop​(Duration timeout)
        Stop computations when input streams are empty. The timeout is applied for each computation, the total duration can be up to nb computations * timeout

        Returns true if computations are stopped during the timeout delay.

      • shutdown

        void shutdown()
        Shutdown immediately.
      • getLowWatermark

        long getLowWatermark​(String computationName)
        Returns the low watermark for the computation. Any message with an offset below the low watermark has been processed by this computation and its ancestors. The returned watermark is local to this processing node, if the computation is distributed the global low watermark is the minimum of all nodes low watermark.
      • getLowWatermark

        long getLowWatermark()
        Returns the low watermark for all the computations of the topology. Any message with an offset below the low watermark has been processed. The returned watermark is local to this processing node.
      • getLatency

        Latency getLatency​(String computationName)
        Returns the latency for a computation. This works also for distributed computations.
        Since:
        10.1
      • isDone

        boolean isDone​(long timestamp)
        Returns true if all messages with a lower timestamp has been processed by the topology.
      • waitForAssignments

        boolean waitForAssignments​(Duration timeout)
                            throws InterruptedException
        Wait for the computations to have assigned partitions ready to process records. The processor must be started. This is useful for writing unit test.

        Returns true if all computations have assigned partitions during the timeout delay.

        Throws:
        InterruptedException
      • isTerminated

        boolean isTerminated()
        True if there is no active processing threads.
        Since:
        10.1
      • toJson

        String toJson​(Map<String,​String> meta)
        Returns a JSON representation of the processor, this includes the list of streams and computations with their settings along with the topology.
        Since:
        11.5