Java Walkthrough

Like all great introductions in software, this walkthrough will start at the beginning: saying hello. The application will run a simple function that accepts a request and responds with a greeting. It will not attempt to cover all the complexities of application development, but instead focus on building a stateful function — which is where you will implement your business logic.

A Basic Hello

Greeting actions are triggered by consuming, routing and passing messages that are defined using ProtoBuf.

  1. syntax = "proto3";
  2. message GreetRequest {
  3. string who = 1;
  4. }
  5. message GreetResponse {
  6. string who = 1;
  7. string greeting = 2;
  8. }

Under the hood, messages are processed using stateful functions, by definition any class that implements the StatefulFunction interface.

  1. package org.apache.flink.statefun.examples.greeter;
  2. import org.apache.flink.statefun.sdk.Context;
  3. import org.apache.flink.statefun.sdk.StatefulFunction;
  4. public final class GreetFunction implements StatefulFunction {
  5. @Override
  6. public void invoke(Context context, Object input) {
  7. GreetRequest greetMessage = (GreetRequest) input;
  8. GreetResponse response = GreetResponse.newBuilder()
  9. .setWho(greetMessage.getWho())
  10. .setGreeting("Hello " + greetMessage.getWho())
  11. .build();
  12. context.send(GreetingConstants.GREETING_EGRESS_ID, response);
  13. }
  14. }

This function takes in a request and sends a response to an external system (or egress). While this is nice, it does not show off the real power of stateful functions: handling state.

A Stateful Hello

Suppose you want to generate a personalized response for each user depending on how many times they have sent a request.

  1. private static String greetText(String name, int seen) {
  2. switch (seen) {
  3. case 0:
  4. return String.format("Hello %s !", name);
  5. case 1:
  6. return String.format("Hello again %s !", name);
  7. case 2:
  8. return String.format("Third times the charm! %s!", name);
  9. case 3:
  10. return String.format("Happy to see you once again %s !", name);
  11. default:
  12. return String.format("Hello at the %d-th time %s", seen + 1, name);
  13. }

Routing Messages

To send a user a personalized greeting, the system needs to keep track of how many times it has seen each user so far. Speaking in general terms, the simplest solution would be to create one function for every user and independently track the number of times they have been seen. Using most frameworks, this would be prohibitively expensive. However, stateful functions are virtual and do not consume any CPU or memory when not actively being invoked. That means your application can create as many functions as necessary — in this case, users — without worrying about resource consumption.

Whenever data is consumed from an external system (or ingress), it is routed to a specific function based on a given function type and identifier. The function type represents the Class of function to be invoked, such as the Greeter function, while the identifier (GreetRequest#getWho) scopes the call to a specific virtual instance based on some key.

  1. package org.apache.flink.statefun.examples.greeter;
  2. import org.apache.flink.statefun.examples.kafka.generated.GreetRequest;
  3. import org.apache.flink.statefun.sdk.io.Router;
  4. final class GreetRouter implements Router<GreetRequest> {
  5. @Override
  6. public void route(GreetRequest message, Downstream<GreetRequest> downstream) {
  7. downstream.forward(GreetingConstants.GREETER_FUNCTION_TYPE, message.getWho(), message);
  8. }
  9. }

So, if a message for a user named John comes in, it will be shipped to John’s dedicated Greeter function. In case there is a following message for a user named Jane, a new instance of the Greeter function will be spawned.

Persistence

Persisted value is a special data type that enables stateful functions to maintain fault-tolerant state scoped to their identifiers, so that each instance of a function can track state independently. To “remember” information across multiple greeting messages, you then need to associate a persisted value field (count) to the Greet function. For each user, functions can now track how many times they have been seen.

  1. package org.apache.flink.statefun.examples.greeter;
  2. import org.apache.flink.statefun.sdk.Context;
  3. import org.apache.flink.statefun.sdk.StatefulFunction;
  4. import org.apache.flink.statefun.sdk.annotations.Persisted;
  5. import org.apache.flink.statefun.sdk.state.PersistedValue;
  6. public final class GreetFunction implements StatefulFunction {
  7. @Persisted
  8. private final PersistedValue<Integer> count = PersistedValue.of("count", Integer.class);
  9. @Override
  10. public void invoke(Context context, Object input) {
  11. GreetRequest greetMessage = (GreetRequest) input;
  12. GreetResponse response = computePersonalizedGreeting(greetMessage);
  13. context.send(GreetingConstants.GREETING_EGRESS_ID, response);
  14. }
  15. private GreetResponse computePersonalizedGreeting(GreetRequest greetMessage) {
  16. final String name = greetMessage.getWho();
  17. final int seen = count.getOrDefault(0);
  18. count.set(seen + 1);
  19. String greeting = greetText(name, seen);
  20. return GreetResponse.newBuilder()
  21. .setWho(name)
  22. .setGreeting(greeting)
  23. .build();
  24. }
  25. }

Each time a message is processed, the function computes a personalized message for that user. It reads and updates the number of times that user has been seen and sends a greeting to the egress.

You can check the full code for the application described in this walkthrough here. In particular, take a look at the module GreetingModule, which is the main entry point for the full application, to see how everything gets tied together. You can run this example locally using the provided Docker setup.

  1. $ docker-compose build
  2. $ docker-compose up

Then, send some messages to the topic “names”, and observe what comes out of “greetings”.

  1. $ docker-compose exec kafka-broker kafka-console-producer.sh \
  2. --broker-list localhost:9092 \
  3. --topic names
  4. docker-compose exec kafka-broker kafka-console-consumer.sh \
  5. --bootstrap-server localhost:9092 \
  6. --isolation-level read_committed \
  7. --from-beginning \
  8. --topic greetings

Java Walkthrough - 图1

Want To Go Further?

This Greeter never forgets a user. Try and modify the function so that it will reset the count for any user that spends more than 60 seconds without interacting with the system.