Skip to content

Introduce Stream API #23

Description

@vesense

Currently, there are Producer API, Consumer API for basic IO operations on message queues.
For many scenes, users like to process messages in Streams.

The Stream API is lightweight, not as same as the distributed streaming systems like Apache Storm, Flink, Spark-Streaming. And its implementation can be embeded in any java applications.

It provides the most common stream operations like filter,flatMap,map,reduce,groupBy,join,count,max,min,window, etc.

It's different from the existing StreamingConsumer, which is a low level Consumer API for processing consumer positions, etc. (IMHO, the class name of StreamingConsumer should be renamed to a more reasonable one, the current might be confused for users.)

Activity

  1. bsideup commented on Mar 19, 2018

    @bsideup

    @vesense I created #20 for it. If OpenMessaging provides Reactive API then such operations can be plugged in by using RxJava, Akka Streams, Reactor or any other reactive framework

  2. vesense commented on Mar 20, 2018

    @vesense
    Author

    @bsideup Thanks for your comment.

    The Stream API here is non-reactive, not as same as Reactive Streams.

    The reactive-streams provides Publisher, Subscriber, Subscription, Processor APIs. I think we might consider providing reactive APIs in an individual module. These APIs can be used as a high level abstraction to process messages in the message queues.

  3. bsideup commented on Mar 20, 2018

    @bsideup

    @vesense but why not Reactive? Does it make sense to provide reactive and non-reactive streams?

  4. vesense commented on Mar 22, 2018

    @vesense
    Author

    @bsideup
    The Stream API here, like Spark's DStream, Flink's DataStream or Kafka's KStream, is just stream abstraction including basic stream operations(filter,flatMap,map,reduce,groupBy,join,count,max,min,window, etc.). This is the different point from reactive-streams's common reactive APIs(Publisher, Subscriber, Subscription, Processor). reactive-streams provides these operations depending on 'RxJava, Akka Streams, Reactor or any other reactive framework'. And another point is that it's easy and convenient to translate sql to Stream related operations when we consider supporting sql(e.g. flink-sql, kafka-ksql).

  5. vongosling commented on May 17, 2018

    @vongosling
    Collaborator

    @vesense how do you think about if we want to support the Stream API base on the existing atomic APIs :-)

  6. vesense commented on May 19, 2018

    @vesense
    Author

    @vongosling Sorry for the delay because of my busy work. I will write a google doc for the Stream API design after completing the RocketMQ-Beam integration. Thanks for your patience.

  7. added this to the 1.1.0 milestone on Sep 4, 2018
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    No type

    Projects

    No projects

      Milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions