Hi All, I am just getting started on FLINK statefu...
# random
g
Hi All, I am just getting started on FLINK stateful functions to develop a high throughput risk orchestrator engine. I am planning to develop Flink App on AWS Kinesis Data Analytics using AWS Lambda (FaaS). While I could develop the stateful functions that can be deployed on AWS lambda, I am unable to get reference /sample code to call these Stateful functions hosted on AWS Lambda from my flink app. A sample code that I got looks like below, but I am unable to find required maven dependency for:
import org.apache.flink.streaming.connectors.aws.lambda.LambdaFunction
and
import org.apache.flink.streaming.connectors.aws.lambda.AWSLambdaInvoke
I am using flink-connector-kinesis-1.15.2.jar in dependency to resolve Flink connectors. Any pointers to resolve above or alternate way of invoking StateFns hosted on AWS lambda using Java code?
Copy code
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.connector.kinesis.sink.KinesisStreamsSink;
import org.apache.flink.kinesis.shaded.com.amazonaws.auth.DefaultAWSCredentialsProviderChain;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kinesis.FlinkKinesisConsumer;
import org.apache.flink.streaming.connectors.kinesis.config.AWSConfigConstants;
import org.apache.flink.streaming.connectors.kinesis.config.ConsumerConfigConstants;
import org.apache.flink.streaming.connectors.aws.lambda.LambdaFunction;
import org.apache.flink.streaming.connectors.aws.lambda.AWSLambdaInvoke;
import org.apache.flink.streaming.connectors.aws.lambda.LambdaInvokeInput;

import java.util.Properties;

/**
 * A basic Kinesis Data Analytics for Java application with Kinesis data
 * streams as source and sink.
 */
public class StateFnsKinesisApp {
    private static final String region = "us-west-2";
    private static final String inputStreamName = "ExampleInputStream";
    private static final String outputStreamName = "ExampleOutputStream";

    private static DataStream<String> createSourceFromStaticConfig(StreamExecutionEnvironment env) {
        Properties inputProperties = new Properties();
        inputProperties.setProperty(ConsumerConfigConstants.AWS_REGION, region);
        inputProperties.setProperty(ConsumerConfigConstants.STREAM_INITIAL_POSITION, "LATEST");

        return env.addSource(new FlinkKinesisConsumer<>(inputStreamName, new SimpleStringSchema(), inputProperties));
    }

    private static KinesisStreamsSink<String> createSinkFromStaticConfig() {
        Properties outputProperties = new Properties();
        outputProperties.setProperty(AWSConfigConstants.AWS_REGION, region);

        return KinesisStreamsSink.<String>builder()
                .setKinesisClientProperties(outputProperties)
                .setSerializationSchema(new SimpleStringSchema())
                .setStreamName(outputProperties.getProperty("OUTPUT_STREAM", "ExampleOutputStream"))
                .setPartitionKeyGenerator(element -> String.valueOf(element.hashCode()))
                .build();
    }

    public static void main(String[] args) throws Exception {
        // Set up the streaming execution environment
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStream<String> inputStream = createSourceFromStaticConfig(env);

        DataStream<String> outputStream = inputStream.map(new MapFunction<String, String>() {
            @Override
            public String map(String value) throws Exception {
                // Process input and return output
                return value;
            }
        });

        DataStream<String> output = outputStream
                .map(new LambdaFunction<String, String>() {
                    @Override
                    public LambdaInvokeInput<String> map(String value) {
                        return LambdaInvokeInput.fromJsonString(value);
                    }
                })
                .addSink(new AWSLambdaInvoke<String>("arn:aws:lambda:us-west-2:123456789012:function:my-stateful-function")
                        .withCredentialsProvider(new DefaultAWSCredentialsProviderChain()));

        env.execute("Flink Stateful function application.");
    }
}
d
FYI, this belongs in #C03G7LJTS2G, not #C03JKTFFX0S. Have you looked at https://flink.apache.org/news/2020/10/13/stateful-serverless-internals.html, which points to https://github.com/tzulitai/statefun-aws-demo/blob/master/app/shopping_cart.py? Might be out-of-date, but it’s a full demo of statefun using AWS lambda.
g
Yes David. This is python code. I am looking for equivalent java library. Problem is I am not able to resolve correct maven dependencies in java for
import org.apache.flink.streaming.connectors.aws.lambda.LambdaFunction
and
import org.apache.flink.streaming.connectors.aws.lambda.AWSLambdaInvoke
j
This does not look like Stateful Function code which generally uses the
org.apache.flink.statefun.flink.core
package. You can see some examples here, and Lambda now has functionality where you can call it via a function URL. Hope this helps.