I/O Module

Stateful Functions’ I/O modules allow functions to receive and send messages to external systems. Based on the concept of Ingress (input) and Egress (output) points, and built on top of the Apache Flink® connector ecosystem, I/O modules enable functions to interact with the outside world through the style of message passing.

Ingress

An Ingress is an input point where data is consumed from an external system and forwarded to zero or more functions. It is defined via an IngressIdentifier and an IngressSpec.

An ingress identifier, similar to a function type, uniquely identifies an ingress by specifying its input type, a namespace, and a name.

The spec defines the details of how to connect to the external system, which is specific to each individual I/O module. Each identifier-spec pair is bound to the system inside an stateful function module.

  1. version: "1.0"
  2. module:
  3. meta:
  4. type: remote
  5. spec:
  6. ingresses:
  7. - ingress:
  8. meta:
  9. id: example/user-ingress
  10. type: # ingress type
  11. spec: # ingress specific configurations
  1. package org.apache.flink.statefun.docs.io.ingress;
  2. import org.apache.flink.statefun.docs.models.User;
  3. import org.apache.flink.statefun.sdk.io.IngressIdentifier;
  4. public final class Identifiers {
  5. public static final IngressIdentifier<User> INGRESS =
  6. new IngressIdentifier<>(User.class, "example", "user-ingress");
  7. }
  1. package org.apache.flink.statefun.docs.io.ingress;
  2. import java.util.Map;
  3. import org.apache.flink.statefun.docs.io.MissingImplementationException;
  4. import org.apache.flink.statefun.docs.models.User;
  5. import org.apache.flink.statefun.sdk.io.IngressIdentifier;
  6. import org.apache.flink.statefun.sdk.io.IngressSpec;
  7. import org.apache.flink.statefun.sdk.spi.StatefulFunctionModule;
  8. public class ModuleWithIngress implements StatefulFunctionModule {
  9. @Override
  10. public void configure(Map<String, String> globalConfiguration, Binder binder) {
  11. IngressSpec<User> spec = createIngress(Identifiers.INGRESS);
  12. binder.bindIngress(spec);
  13. }
  14. private IngressSpec<User> createIngress(IngressIdentifier<User> identifier) {
  15. throw new MissingImplementationException("Replace with your specific ingress");
  16. }
  17. }

Router

A router is a stateless operator that takes each record from an ingress and routes it to zero or more functions. Routers are bound to the system via a stateful function module, and unlike other components, an ingress may have any number of routers.

When defined in yaml, routers are defined by a list of function types. The id component of the address is pulled from the key associated with each record in its underlying source implementation.

  1. targets:
  2. - example-namespace/my-function-1
  3. - example-namespace/my-function-2
  1. package org.apache.flink.statefun.docs.io.ingress;
  2. import org.apache.flink.statefun.docs.FnUser;
  3. import org.apache.flink.statefun.docs.models.User;
  4. import org.apache.flink.statefun.sdk.io.Router;
  5. public class UserRouter implements Router<User> {
  6. @Override
  7. public void route(User message, Downstream<User> downstream) {
  8. downstream.forward(FnUser.TYPE, message.getUserId(), message);
  9. }
  10. }
  1. package org.apache.flink.statefun.docs.io.ingress;
  2. import java.util.Map;
  3. import org.apache.flink.statefun.docs.models.User;
  4. import org.apache.flink.statefun.sdk.io.IngressIdentifier;
  5. import org.apache.flink.statefun.sdk.io.IngressSpec;
  6. import org.apache.flink.statefun.sdk.io.Router;
  7. import org.apache.flink.statefun.sdk.spi.StatefulFunctionModule;
  8. public class ModuleWithRouter implements StatefulFunctionModule {
  9. @Override
  10. public void configure(Map<String, String> globalConfiguration, Binder binder) {
  11. IngressSpec<User> spec = createIngressSpec(Identifiers.INGRESS);
  12. Router<User> router = new UserRouter();
  13. binder.bindIngress(spec);
  14. binder.bindIngressRouter(Identifiers.INGRESS, router);
  15. }
  16. private IngressSpec<User> createIngressSpec(IngressIdentifier<User> identifier) {
  17. throw new MissingImplementationException("Replace with your specific ingress");
  18. }
  19. }

Egress

Egress is the opposite of ingress; it is a point that takes messages and writes them to external systems. Each egress is defined using two components, an EgressIdentifier and an EgressSpec.

An egress identifier uniquely identifies an egress based on a namespace, name, and producing type. An egress spec defines the details of how to connect to the external system, the details are specific to each individual I/O module. Each identifier-spec pair are bound to the system inside a stateful function module.

  1. version: "1.0"
  2. module:
  3. meta:
  4. type: remote
  5. spec:
  6. egresses:
  7. - egress:
  8. meta:
  9. id: example/user-egress
  10. type: # egress type
  11. spec: # egress specific configurations
  1. package org.apache.flink.statefun.docs.io.egress;
  2. import org.apache.flink.statefun.docs.models.User;
  3. import org.apache.flink.statefun.sdk.io.EgressIdentifier;
  4. public final class Identifiers {
  5. public static final EgressIdentifier<User> EGRESS =
  6. new EgressIdentifier<>("example", "egress", User.class);
  7. }
  1. package org.apache.flink.statefun.docs.io.egress;
  2. import java.util.Map;
  3. import org.apache.flink.statefun.docs.io.MissingImplementationException;
  4. import org.apache.flink.statefun.docs.models.User;
  5. import org.apache.flink.statefun.sdk.io.EgressIdentifier;
  6. import org.apache.flink.statefun.sdk.io.EgressSpec;
  7. import org.apache.flink.statefun.sdk.spi.StatefulFunctionModule;
  8. public class ModuleWithEgress implements StatefulFunctionModule {
  9. @Override
  10. public void configure(Map<String, String> globalConfiguration, Binder binder) {
  11. EgressSpec<User> spec = createEgress(Identifiers.EGRESS);
  12. binder.bindEgress(spec);
  13. }
  14. public EgressSpec<User> createEgress(EgressIdentifier<User> identifier) {
  15. throw new MissingImplementationException("Replace with your specific egress");
  16. }
  17. }

Stateful functions may then message an egress the same way they message another function, passing the egress identifier as function type.

  1. package org.apache.flink.statefun.docs.io.egress;
  2. import org.apache.flink.statefun.docs.models.User;
  3. import org.apache.flink.statefun.sdk.Context;
  4. import org.apache.flink.statefun.sdk.StatefulFunction;
  5. /** A simple function that outputs messages to an egress. */
  6. public class FnOutputting implements StatefulFunction {
  7. @Override
  8. public void invoke(Context context, Object input) {
  9. context.send(Identifiers.EGRESS, new User());
  10. }
  11. }