Flink custom parallel data source

2023-10-30 Thread Kamal Mittal via user
Hello Community, I need to have a custom parallel data source (Flink ParallelSourceFunction) for fetching data based on some custom logic. In this source function, opening multiple threads via java thread pool to distribute work further. These threads share Flink provided 'SourceContext' and co

Re:Query on Flink SQL create DDL primary key for nested field

2023-10-30 Thread Xuyang
Hi, Flink SQL doesn't support a inline field in struct type as pk. You can try to raise an issue about this feature in community[1]. For a quick solution, you can try to transform it by DataStream API first by extracting the 'id' and then convert it to Table API to use SQL. [1] https://issu

Best practice way to conditionally discard a window and not serialize the results

2023-10-30 Thread Mark Petronic
I am reading stats from Kinesis, deserializing them into a stat POJO and then doing something like this using an aggregated window with no defined processWindow function: timestampedStats .keyBy(v -> v.groupKey()) .window(TumblingEventTimeWindows.of(Time.seconds(appCfg.get

Re: [DISCUSS][FLINK-33240] Document deprecated options as well

2023-10-30 Thread Matthias Pohl via user
Thanks for your proposal, Zhanghao Chen. I think it adds more transparency to the configuration documentation. +1 from my side on the proposal On Wed, Oct 11, 2023 at 2:09 PM Zhanghao Chen wrote: > Hi Flink users and developers, > > Currently, Flink won't generate doc for the deprecated options

Re: Metrics with labels

2023-10-30 Thread Lars Skjærven
Registering the counter is fine, e.g. in `open()`: lazy val responseCounter: Counter = getRuntimeContext .getMetricGroup .addGroup("response_code") .counter("myResponseCounter") then, in def asyncInvoke(), I can still only do responseCounter.inc(), but what I want is responseCou

Clear the State Backends in Flink

2023-10-30 Thread arjun s
Hi team, I'm interested in understanding if there is a method available for clearing the State Backends in Flink. If so, could you please provide guidance on how to accomplish this particular use case? Thanks and regards, Arjun S

Monitoring File Processing Progress in Flink Jobs

2023-10-30 Thread arjun s
Hi team, I'm also interested in finding out if there is Java code available to determine the extent to which a Flink job has processed files within a directory. Additionally, I'm curious about where the details of the processed files are stored within Flink. Thanks and regards, Arjun S

Re: Enhancing File Processing and Kafka Integration with Flink Jobs

2023-10-30 Thread arjun s
Hi team, I'm also interested in finding out if there is Java code available to determine the extent to which a Flink job has processed files within a directory. Additionally, I'm curious about where the details of the processed files are stored within Flink. Thanks and regards, Arjun S On Mon, 30

Checkpoints are not triggering when S3 is unavailable

2023-10-30 Thread Evgeniy Lyutikov
Hi team! I came across strange behavior in Flink 1.17.1. If during the build of a checkpoint the s3 storage becomes unavailable, then the current checkpoint expired by timeout and new ones are not triggered. The triggering for new checkpoints is resumed only after s3 is restored and this can be

[ANNOUNCE] Apache Flink Kubernetes Operator 1.6.1 released

2023-10-30 Thread Rui Fan
The Apache Flink community is very happy to announce the release of Apache Flink Kubernetes Operator 1.6.1. Please check out the release blog post for an overview of the release: https://flink.apache.org/2023/10/27/apache-flink-kubernetes-operator-1.6.1-release-announcement/ The release is availa

Query on Flink SQL create DDL primary key for nested field

2023-10-30 Thread elakiya udhayanan
Hi team, I have a Kafka topic named employee which uses confluent avro schema and will emit the payload as below: { "employee": { "id": "123456", "name": "sampleName" } } I am using the upsert-kafka connector to consume the events from the above Kafka topic as below using the Flink SQL DDL statem