Gaurav Gupta
01/17/2023, 2:53 PMimport 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?
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.");
}
}David Anderson
01/17/2023, 3:15 PMGaurav Gupta
01/17/2023, 3:17 PMimport org.apache.flink.streaming.connectors.aws.lambda.LambdaFunction and import org.apache.flink.streaming.connectors.aws.lambda.AWSLambdaInvokeJeremy Ber
01/20/2023, 2:25 PMorg.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.