Skip to content
CAI
Software that uses CAICheck a score

apache/flink

52.6

Adequate · 24 September 2026

1451.4k

lines of production code

Java

with Scala

4

measurements over time

CAI band scale
CAI trend line
CAI lens gauges

What this system is

This system is Apache Flink, a distributed stream processing framework that manages the lifecycle of data pipelines through execution, state management, and fault tolerance. It provides a unified programming model supporting both streaming and batch workloads via the DataStream and Table/SQL APIs, while handling complex state operations through various backends like RocksDB and ForSt. The platform integrates with diverse external systems, including cloud storage services, Hadoop ecosystems, and machine learning models, and exposes its capabilities through REST APIs, JDBC drivers, and a web dashboard for monitoring and management.

How it got here

2010–2018 — Infrastructure modernization and API expansion

78 changes.

This period focused on modernizing the project's development infrastructure through comprehensive dependency convergence, JUnit 5 migration, and the introduction of Azure Pipelines CI. It also involved significant architectural shifts, including the restructuring of the binary distribution, the implementation of new metrics reporting interfaces, and the expansion of Hadoop and S3 compatibility layers. Concurrently, the work established robust testing frameworks and introduced new features for the DataStream and Table APIs, such as Avro and Parquet bulk I/O support.

2019–2020 — SQL parser overhaul and connector modernization

62 changes.

This period focused on a major overhaul of the SQL parser to support extensive new syntax and data types, alongside a comprehensive restructuring of the Table API and SQL Client architecture. Significant work was done to modernize connectors, introducing new FileSink capabilities, Azure storage support, and generic async sink frameworks, while expanding test coverage for classloader isolation and state migration.

2021–2023 — RPC migration and SQL Gateway introduction

62 changes.

This period focused on migrating the RPC layer from Akka to Apache Pekko and introducing the SQL Gateway REST API for external session management. It also established robust architectural testing via ArchUnit and expanded state management capabilities with the ChangelogStateBackend and native S3 support.

2024–2026 — DataStream V2 and AI integration

15 changes.

This period focused on introducing the DataStream V2 API, establishing new core interfaces for state and execution modes, and implementing the associated runtime infrastructure. It also expanded Flink's capabilities by adding native S3 support, OpenTelemetry observability, and integrations with OpenAI and Triton AI models, alongside significant test coverage and utility improvements.

Features

Add Aliyun OSS file system without Hadoop dependencies

Introduces a new Aliyun OSS filesystem implementation for Flink that operates independently of Hadoop dependencies. This includes the core filesystem classes (FlinkOSSFileSystem, OSSAccessor, OSSFileSystemFactory), a recoverable writer for checkpointing and streaming (OSSRecoverableWriter, OSSRecoverableFsDataOutputStream, OSSRecoverableMultipartUpload), and the necessary service provider configuration. The implementation supports multipart uploads with configurable part sizes and concurrent upload limits, ensuring data durability through recoverable state serialization.

flink-filesystems/flink-oss-fs-hadoop · high confidence

Add Azure Blob Storage and ADLS Gen2 filesystem support

This change introduces the \flink-azure-fs-hadoop\ module, enabling Flink to read from and write to Azure Blob Storage (via \wasb\/\wasbs\ schemes) and Azure Data Lake Storage Gen2 (via \abfs\/\abfss\ schemes). It provides the necessary \FileSystemFactory\ implementations, a recoverable writer for streaming sinks, and an \EnvironmentVariableKeyProvider\ to retrieve storage keys from the \AZURE\_STORAGE\_KEY\ environment variable. The module also includes integration tests and architecture tests to validate the new functionality.

flink-filesystems/flink-azure-fs-hadoop · high confidence

Add Confluent Schema Registry Avro format with SSL and auth support

Users can now configure the 'avro-confluent' Table format to read from and write to Kafka topics using the Confluent Schema Registry. This change introduces a new format factory that handles schema registration and lookup, and adds configuration options for SSL keystore/truststore locations and passwords, as well as basic and bearer authentication credentials. The format supports both generic and specific Avro records and allows forwarding additional registry properties via a generic 'properties' map.

flink-formats/flink-avro-confluent-registry · high confidence

Add Hadoop mapred Input/Output Format wrappers for the Java API

This change introduces new wrapper classes (HadoopInputFormat, HadoopOutputFormat, and their bases) in the flink-hadoop-compatibility module, enabling users to read from and write to Hadoop mapred InputFormats and OutputFormats directly using Flink's Java DataSet API. These classes handle the necessary bridging between Flink's execution model and Hadoop's mapred interfaces, including configuration, split assignment, and record reading/writing.

flink-connectors/flink-hadoop-compatibility/src/main/java/org/apache/flink/api/java/hadoop/mapred · high confidence

Add Java API Hadoop MapReduce Input/Output Format support

This change introduces the Java API-specific implementations for reading from and writing to Hadoop MapReduce formats. It adds the \HadoopInputFormat\ and \HadoopOutputFormat\ classes, along with their shared base classes (\HadoopInputFormatBase\, \HadoopOutputFormatBase\) and supporting utilities (\HadoopUtils\, \HadoopInputSplit\). These components enable Flink jobs using the Java API to directly consume Hadoop InputFormats and produce Hadoop OutputFormats, handling necessary configuration merging, split wrapping, and serialization.

flink-connectors/flink-hadoop-compatibility/src/main/java/org/apache/flink/api/java/hadoop/mapreduce · high confidence

Add ORC format support without Hive dependencies

This change introduces a new \flink-orc-nohive\ module that provides ORC read and write capabilities without requiring the Hive runtime. It includes a \BulkWriter.Factory\ for writing ORC files, a columnar row input format and split reader for reading ORC files, and a set of vector adapters that map ORC's column vectors to Flink's internal columnar data structures. This allows users to process ORC data in Flink without the overhead or complexity of the full Hive dependency.

flink-formats/flink-orc-nohive · high confidence

Add SLF4J-based reporters for metrics, traces, and events

The flink-metrics-slf4j module now provides three new reporters that export Flink telemetry to SLF4J loggers: a Slf4jReporter for metrics (counters, gauges, meters, histograms), a Slf4jTraceReporter for distributed traces (spans), and a Slf4jEventReporter for application events. Each reporter is registered via META-INF/services factories so they can be enabled through standard reporter configuration, and the reporter includes logic to skip logging when no metrics are registered and to safely handle concurrent modifications during reporting.

flink-metrics/flink-metrics-slf4j · high confidence

Add SequenceFile support to the new FileSink API

This change introduces a new SequenceFile format implementation for Flink's modern FileSink. It adds a \SequenceFileWriterFactory\ and \SequenceFileWriter\ that implement the \BulkWriter\ interface, allowing users to write SequenceFile outputs with configurable Hadoop compression codecs. The module also includes a \SerializableHadoopConfiguration\ wrapper to safely serialize Hadoop configurations across the cluster, along with integration tests verifying the end-to-end writing and reading of SequenceFiles via the new FileSink API.

flink-formats/flink-sequence-file · high confidence

Add StatsD metrics reporter

Introduces a new StatsDReporter that sends Flink metrics (counters, gauges, histograms, and meters) to a StatsD-compatible backend via UDP. The reporter is automatically discovered via the Java Service Provider Interface (SPI) using StatsDReporterFactory, and users can configure the target host and port through the reporter configuration. It also handles invalid characters in metric names and properly reports negative values.

flink-metrics/flink-metrics-statsd · high confidence

Add TPC-DS end-to-end test suite

This change introduces a new end-to-end test module for the TPC-DS benchmark. It includes a test program that executes 99 TPC-DS queries against a batch-mode TableEnvironment, registering 24 source tables with their schemas and optional statistics. The suite also provides utilities to format TPC-DS answer sets and compare Flink's query results against the standard expected outputs, ensuring strict compliance with the TPC-DS specification.

flink-end-to-end-tests/flink-tpcds-test · high confidence

Add TPC-H benchmark data generation and result validation tools

The flink-tpch-test module now includes utilities to generate TPC-H test data and validate query results. TpchDataGenerator creates standard TPC-H tables, SQL queries, and expected output files based on a scale factor, while TpchResultComparator verifies actual query outputs against expected results, implementing the rounding and comparison logic specified in the TPC-H standard specification v2.18.0.

flink-end-to-end-tests/flink-tpch-test · high confidence

Add WritableTypeInfo for Hadoop Writable types

The flink-hadoop-compatibility module now includes WritableTypeInfo, a new TypeInformation implementation for data types extending Hadoop's Writable interface. This enables Flink to properly handle serialization and comparison for Hadoop Writable types within the unified connector module.

flink-connectors/flink-hadoop-compatibility/src/main/java/org/apache/flink/api/java/typeutils · high confidence

Add common data models and sources for the DataStream API walkthrough

The \flink-walkthrough-common\ module now includes the \Alert\ and \Transaction\ entity classes, along with a \TransactionSource\ that generates simulated credit card transactions using the \DataGeneratorSource\ API, and sink utilities (\AlertSink\, \LoggerOutputFormat\) to support the new DataStream API fraud detection walkthrough.

flink-walkthroughs/flink-walkthrough-common · high confidence

Add heavy deployment stress test for inflated operator state

A new end-to-end stress test program has been added to validate Flink's behavior under heavy deployment descriptors. The test creates a source operator that initializes a large amount of union operator state (configurable via \num\_list\_states\_per\_op\ and \num\_partitions\_per\_list\_state\), effectively inflating the metadata size during deployment. This allows users to verify system stability and performance when handling tasks with significantly expanded state metadata.

flink-end-to-end-tests/flink-heavy-deployment-stress-test · high confidence

CEP library introduces new processing API and configurable SharedBuffer caching

The Flink CEP library adds a new \PatternProcessFunction\ as the preferred way to process detected pattern sequences, providing access to time features and side outputs, while introducing \CEPCacheOptions\ to allow users to configure the maximum number of event and entry slots in the \SharedBuffer\ cache (defaulting to 1024 slots each) and the interval for logging cache statistics. This change also includes the introduction of \EventComparator\ for sorting events with equal timestamps and \RichPatternSelectFunction\/\RichPatternFlatSelectFunction\ variants that provide access to the \RuntimeContext\.

flink-libraries/flink-cep · high confidence

CSV format now supports projection pushdown and statistics reporting for the File System connector

The CSV format module introduces a new \CsvFileFormatFactory\ that implements \BulkReaderFormatFactory\ and \BulkWriterFormatFactory\ for the File System connector. This enables projection pushdown, allowing the connector to read only the required columns from CSV files, and adds statistics reporting capabilities to improve query planning. The implementation includes new classes such as \AbstractCsvInputFormat\ for optimized file splitting, \CsvBulkWriter\ for efficient bulk writing, and \CsvReaderFormat\ for stream-based reading, alongside updated \CsvFormatOptions\ to expose configuration options like \ignore-first-line\ and \disable-quote-character\ to the table connector.

flink-formats/flink-csv · high confidence

The connector test framework now includes a new container-based test environment that runs a real Flink cluster (JobManager and TaskManagers) inside Docker containers using Testcontainers. This allows connector tests to execute against a distributed Flink setup rather than an in-process mini-cluster, supporting features like Zookeeper-based high availability, custom checkpoint and recovery paths, and network isolation. The environment is exposed via the new \FlinkContainerTestEnvironment\ class and is managed through \FlinkContainers\, \FlinkContainersSettings\, and \TestcontainersSettings\ to configure cluster topology, base images, and dependencies.

flink-test-utils-parent/flink-connector-test-utils · high confidence

DataStream API V2 introduces experimental ProcessFunction-based programming model

The DataStream API V2 introduces a new experimental programming model centered on ProcessFunction, replacing the previous lambda-based approach. This change provides richer runtime context through new interfaces like RuntimeContext, PartitionedContext, and NonPartitionedContext, exposing job and task metadata, state management, and processing time capabilities. The API also adds built-in support for event-time watermarks via EventTimeExtension, join operations through BuiltinFuncs, and window processing, enabling more complex stream processing patterns with explicit state and timer management.

flink-datastream-api · high confidence

DataStream V2 API implementation

This release introduces the implementation for the new DataStream V2 API, providing the core execution environment and runtime context infrastructure. The \ExecutionContextEnvironment\ and \ExecutionEnvironmentImpl\ classes establish the entry point for building and executing DataStream V2 jobs, supporting both local and cluster execution modes. The new context hierarchy (\DefaultRuntimeContext\, \DefaultPartitionedContext\, \DefaultNonPartitionedContext\) exposes essential runtime information such as \JobInfo\, \TaskInfo\, and \MetricGroup\ to user functions. Additionally, the \DefaultStateManager\ implements the V2 state access API, enabling users to interact with Value, List, Map, Reducing, and Aggregating states through the new \StateDeclaration\ interface, while \DefaultProcessingTimeManager\ and \DefaultWatermarkManager\ provide the underlying mechanisms for timer and watermark management.

flink-datastream · high confidence

Google Cloud Storage filesystem now supports RecoverableWriter

The Google Cloud Storage (GCS) filesystem connector now supports the RecoverableWriter interface, enabling more robust and idempotent writes for Flink's FileSink. This change introduces a new \GSFileSystem\ implementation that wraps the Hadoop GCS connector and provides a \GSRecoverableWriter\ backed by a new abstraction layer (\GSBlobStorage\). Users benefit from improved reliability through configurable retry settings (max attempts, timeouts), HTTP connection/read timeouts, and optional entropy injection in temporary object names to mitigate GCS hotspotting. Additionally, the connector allows specifying a separate bucket for temporary files and supports overriding the GCS root URL via Hadoop configuration.

flink-filesystems/flink-gs-fs-hadoop · high confidence

Graphite reporter now supports configurable TCP/UDP protocol

The Graphite metrics reporter in flink-metrics-graphite now allows users to select the transport protocol (TCP or UDP) via the new 'protocol' configuration property, defaulting to TCP if unspecified. This change introduces a new GraphiteReporter class that wraps the Dropwizard reporter and a GraphiteReporterFactory registered via Java SPI to enable service-loader-based instantiation. The implementation also bundles the Dropwizard metrics-graphite library (version 3.2.6) and includes updated NOTICE/LICENCE files reflecting the 2014-2026 copyright period.

flink-metrics/flink-metrics-graphite · high confidence

The Flink SQL JDBC driver is now available in the flink-table module, enabling users to connect to Flink SQL Gateway via standard JDBC. This new driver supports connection establishment, catalog and schema selection via URI parsing, and basic statement execution. It includes implementations for core JDBC interfaces such as Connection, Statement, ResultSet, and DatabaseMetaData, along with type mapping for Flink logical types to SQL types. The driver is designed to be compatible with standard JDBC tools like Beeline and SQLLine.

flink-table/flink-sql-jdbc-driver · high confidence

Initial repository configuration and developer documentation

The repository now includes foundational configuration files and documentation to standardize the development environment. This adds an \.asf.yaml\ file to configure GitHub repository settings (such as merge strategies and notification channels), an \.editorconfig\ to enforce consistent code formatting across IDEs, and a \.gitignore\ to exclude build artifacts and IDE-specific files. Additionally, it introduces \AGENTS.md\ with guidelines for AI coding agents, \DEVELOPMENT.md\ with detailed instructions for setting up IntelliJ IDEA and other IDEs, and \azure-pipelines.yml\ to define the continuous integration build pipeline.

(repo-wide) · high confidence

Introduce ChangelogStateBackend for state change logging

The ChangelogStateBackend is introduced as a new state backend that wraps an existing delegated backend (such as Heap or RocksDB) to automatically log all state mutations. This enables incremental checkpointing and faster recovery by allowing the system to restore state from the change log rather than re-processing the entire state snapshot. The implementation includes specific changelog wrappers for partitioned state types, including ListState, MapState, AggregatingState, ReducingState, and KeyGroupedInternalPriorityQueue, ensuring that all updates are captured in the changelog stream.

flink-state-backends/flink-statebackend-changelog · high confidence

Introduce DFS-based Changelog State Backend with local caching and metrics

This change adds a new Filesystem-based implementation of the Changelog State Backend (FsStateChangelogStorage) to the flink-dstl-dfs module. It enables storing changelog state in distributed file systems with support for local recovery via a local cache (ChangelogStreamHandleReaderWithCache) that caches remote files to local disk to reduce read latency. The implementation includes a batching upload scheduler (BatchingStateChangeUploadScheduler) to group state changes, configurable retry policies for uploads, and comprehensive metrics (ChangelogStorageMetricGroup) for monitoring upload performance, failures, and queue sizes.

flink-dstl/flink-dstl-dfs · high confidence

Introduce Dropwizard-based metrics reporter

The flink-metrics-dropwizard module now provides a new ScheduledDropwizardReporter that bridges Flink's metric types (Counter, Gauge, Histogram, Meter) to the Dropwizard metrics library. This enables users to export Flink metrics to any backend supported by Dropwizard reporters. The implementation includes wrapper classes to adapt Flink metrics to Dropwizard types and vice versa, ensuring that metric names are sanitized to remove invalid characters. Tests verify correct registration, reporting, and cleanup of metrics through the Dropwizard registry.

flink-metrics/flink-metrics-dropwizard · high confidence

Introduce FLIP-238 DataGeneratorSource

The DataGen connector now includes a new DataGeneratorSource based on FLIP-238, allowing users to generate a bounded stream of N data points in parallel using a custom GeneratorFunction. This source supports rate limiting and deterministic element generation, providing a modern alternative to the legacy FiniteTestSource for testing and scenarios requiring synthetic event streams.

flink-connectors/flink-connector-datagen · high confidence

Introduce ForSt State Backend with comprehensive state management capabilities

This change introduces the ForSt State Backend module, providing a new state storage implementation for Flink. It includes configuration options for tuning background threads, file limits, and compaction styles, as well as implementations for various state types including aggregating, list, and map states. The backend supports asynchronous state access, namespace handling, and integrates with the State V2 API, enabling users to leverage ForSt for improved state management performance and flexibility.

flink-state-backends/flink-statebackend-forst · high confidence

Introduce GPU external resource driver and discovery scripts

Adds a new GPU external resource driver that allows Flink to discover and allocate GPU resources via a configurable discovery script. The implementation includes the GPUDriver and GPUDriverFactory classes, configuration options for the script path and arguments, and a default NVIDIA GPU discovery script (nvidia-gpu-discovery.sh) that supports both non-coordination and coordination modes for resource allocation. Tests verify the driver's behavior and the discovery script's logic.

flink-external-resources/flink-external-resource-gpu · high confidence

Introduce Hadoop-based Recoverable Writer for fault-tolerant streaming sinks

This change adds a new RecoverableWriter implementation for Hadoop file systems (HDFS and ViewFS), enabling the StreamingFileSink to safely resume writes after job failures or checkpoints. The new component includes the HadoopRecoverableWriter, HadoopRecoverableFsDataOutputStream, and supporting classes (HadoopFsRecoverable, HadoopRecoverableSerializer) to manage staging files, track write positions, and handle lease revocation and file truncation during recovery. It also introduces base classes and wrappers for Hadoop file status, block locations, and data streams to integrate Hadoop's native file system operations with Flink's internal filesystem abstraction.

flink-filesystems/flink-hadoop-fs · high confidence

Introduce InfluxDB metrics reporter

Adds a new InfluxDB reporter that allows Flink to export metrics to an InfluxDB database. The reporter supports configurable host, port, database, username, and password, as well as optional retention policies and consistency levels. It maps Flink counters, gauges, histograms, and meters to InfluxDB points with appropriate fields and tags, and includes a factory for service loader registration along with unit and integration tests.

flink-metrics/flink-metrics-influxdb · high confidence

Introduce Java code splitter module for generated code

A new \flink-table-code-splitter\ module has been added to automatically split large generated Java methods into smaller, more manageable pieces. This addresses issues with generated code exceeding JVM method size limits or becoming too complex for the JIT compiler. The module includes ANTLR grammar files (\JavaLexer.g4\, \JavaParser.g4\) for parsing Java code, and several rewriter classes (\FunctionSplitter\, \BlockStatementGrouper\, \DeclarationRewriter\, etc.) that handle splitting functions, grouping block statements, and managing variable declarations to ensure the rewritten code remains valid and efficient.

flink-table/flink-table-code-splitter · high confidence

Introduce Native S3 FileSystem with AWS SDK v2

A new native S3 filesystem module (flink-s3-fs-native) is added, providing a direct implementation of Flink's FileSystem interface using AWS SDK v2 without Hadoop dependencies. It supports s3:// and s3a:// URI schemes for checkpointing and file sinks with exactly-once semantics via multipart uploads. The module includes comprehensive configuration options for credentials, retries, timeouts, server-side encryption (SSE-S3/SSE-KMS), IAM assume role, and bucket-level configuration overrides. It also exposes AWS SDK operation metrics (API call counts, latency, throttles, retries) and supports batch recursive deletes. A helper script (download-crt-jars.sh) is provided to download the aws-crt JAR for enabling CRT-based networking when s3.crt.enabled=true.

flink-filesystems · high confidence

Introduce OpenAI Model Function support for chat and embeddings

Users can now connect Flink to OpenAI models via a new \openai\ ModelProviderFactory. This adds support for both chat completions (with configurable system prompts, temperature, top-p, stop sequences, max tokens, presence penalty, n, seed, and response format) and text embeddings (with configurable dimensions). The implementation includes robust context management, allowing users to set a maximum context size and define how overflows are handled (truncating from the head or tail, or skipping the input, with optional logging). It also features configurable error handling strategies, including retry counts and fallback behaviors, ensuring resilient integration with the OpenAI API.

flink-models/flink-model-openai · high confidence

Introduce Parquet bulk writer and vectorized columnar reader for the File System connector

The Parquet format module now provides a new bulk writer implementation (ParquetBulkWriter) and a vectorized columnar row input format (ParquetColumnarRowInputFormat) to integrate with the Flink File System connector. This enables more efficient, columnar-based reading of Parquet files into Flink's RowData structures and supports writing via the new BulkWriter interface, while also adding Avro-specific readers and writers for Parquet files.

flink-formats/flink-parquet · high confidence

Introduce Pekko RPC system loader with fallback support

The flink-rpc-akka-loader module now provides a PekkoRpcSystemLoader that loads the RPC system from a bundled fat JAR using a dedicated SubmoduleClassLoader, ensuring isolation. A FallbackPekkoRpcSystemLoader is also added to support IDE development by dynamically resolving dependencies via Maven when the fat JAR is not present. Service provider configurations are updated to register these loaders, and integration tests verify the loading behavior under various temporary directory configurations.

flink-rpc/flink-rpc-akka-loader · high confidence

Introduce Presto-based S3 filesystem with custom recursive deletion

This change adds a new Presto-specific S3 filesystem implementation (accessible via the 's3' and 's3p' schemes) that overrides the default recursive deletion behavior to work around a bug in the underlying Presto file system. It includes the necessary factory classes, delegation token providers, and service registrations to enable this functionality, along with mock classes to avoid incompatible dependencies.

flink-filesystems/flink-s3-fs-presto · high confidence

Introduce Protobuf format for Table API streaming

Adds a new 'protobuf' format for the Flink Table API, enabling users to serialize and deserialize RowData to and from Protocol Buffers messages. The implementation includes a format factory (PbFormatFactory) that accepts a required 'message-class-name' option and optional flags for 'ignore-parse-errors', 'read-default-values', and 'write-null-string-literal'. It provides codegen-based deserializers for simple types, rows, arrays, and maps, and explicitly prevents the use of this format with the filesystem connector to avoid silent data corruption.

flink-formats/flink-protobuf · high confidence

Introduce SQL Gateway REST endpoint framework and handlers

The SQL Gateway now exposes a comprehensive REST API for managing sessions, executing SQL statements, and controlling materialized tables. This change introduces the \SqlGatewayRestEndpoint\ and its factory, wiring up handlers for session lifecycle (open, close, configure, heartbeat), statement execution (execute, complete, fetch results), and materialized table operations (refresh, and embedded scheduler workflows for periodic refresh). It also adds a script deployment handler for application mode, utility endpoints for API version and info, and a CLI entry point (\SqlGateway\) with options parsing to start the service.

flink-table/flink-sql-gateway · high confidence

Introduce Standalone Application Cluster entry point and configuration

The flink-container module now includes a new entry point for running Flink jobs in standalone mode via the application paradigm. This change adds the StandaloneApplicationClusterEntryPoint, which parses command-line arguments for job class name, job ID, and JAR files, and initializes the necessary runtime components like the artifact fetch manager and plugin manager. It also introduces the StandaloneApplicationClusterConfiguration and its parser factory to handle these specific parameters, along with corresponding JUnit 5 tests and test infrastructure updates (log4j2, JUnit extensions).

flink-container · high confidence

Introduce Triton Inference Server model function with resilience features

This change adds the Triton Inference Server model provider to Flink, enabling users to invoke Triton models via the \TRITON\ factory identifier. The implementation includes an abstract base class and a concrete inference function that handles HTTP communication, data type mapping, and request serialization. Key capabilities include configurable retry logic with exponential backoff, a health checker that periodically probes the server, and a circuit breaker to fail fast and prevent cascading failures when the Triton server is unhealthy. The module also supports stateful model sequences, gzip compression, and default-value fallbacks for failed inferences.

flink-models/flink-model-triton · high confidence

Introduce WritableSerializer and WritableComparator for Hadoop Writable types

The flink-hadoop-compatibility module now includes new WritableSerializer and WritableComparator classes to handle serialization and comparison of Hadoop Writable types. The WritableSerializer leverages Kryo for copying operations while using native Writable methods for serialization, and includes configuration snapshotting for checkpoint compatibility. The WritableComparator provides type-safe comparison logic for Writable instances, supporting both ascending and descending orderings and normalized key operations where applicable.

flink-connectors/flink-hadoop-compatibility/src/main/java/org/apache/flink/api/java/typeutils/runtime · high confidence

Introduce copy-on-write skip list state map for spillable heap backend

The spillable heap state backend now uses a new copy-on-write skip list implementation to manage state. This change introduces a concurrent, versioned state map that supports efficient snapshotting and on-disk spilling, allowing users to handle larger state sizes with improved concurrency and reduced memory pressure during checkpointing.

flink-state-backends/flink-statebackend-heap-spillable · high confidence

Introduce native S3 filesystem implementation with AWS SDK v2

Adds a new native S3 filesystem implementation built on AWS SDK v2, providing a factory for the s3:// scheme and a compatibility factory for the s3a:// scheme. This implementation introduces bucket-level configuration support (allowing per-bucket overrides for endpoints, credentials, and encryption), configurable bulk copy operations for S3-to-local downloads, and optimized input streams with lazy seeking to reduce IOPS. It also includes batched recursive deletion and exposes operation-level S3 metrics for observability.

flink-filesystems/flink-s3-fs-native · high confidence

Introduce new CLI command-line option classes and pipeline translation utilities

The Flink client now includes a set of dedicated command-line option classes (such as CancelOptions, CheckpointOptions, and ArtifactFetchOptions) to structure CLI argument parsing, alongside new utility classes (ClientUtils, FlinkPipelineTranslationUtil, StreamGraphTranslator) that handle program execution context and translate DataStream Pipelines into JobGraphs. These changes support the new executor-based architecture by centralizing CLI option handling and providing a reflection-based mechanism to discover and apply the correct translator for a given pipeline type.

flink-clients/src/main · high confidence

Introduce new Datadog HTTP metrics reporter

The flink-metrics-datadog module now includes a new HTTP-based metrics reporter (DatadogHttpReporter) that sends metrics to Datadog via its REST API. This reporter supports counters, gauges, meters, and histograms (mapped to gauge sub-metrics), and allows configuration of the API key, proxy settings, target data center (US or EU), and metric tagging. It also introduces a switch to use logical identifiers for metric names and deprecates the legacy 'tags' configuration option in favor of 'scope.variables.additional'.

flink-metrics/flink-metrics-datadog · high confidence

Introduce periodic changelog materialization metrics and manager in state-backend-common

The \flink-state-backends/flink-statebackend-common\ module now includes the \ChangelogMaterializationMetricGroup\ and \PeriodicMaterializationManager\ classes, which expose metrics for tracking started, completed, and failed materialization events along with the duration of the last materialization. This change also adds the corresponding unit tests for the new manager and configures JUnit 5 extensions for the test suite.

flink-state-backends/flink-statebackend-common · high confidence

Introduce pluggable dialect support via new Calcite bridge module

A new internal module, flink-table-calcite-bridge, has been added to support pluggable SQL dialects. This change introduces the CalciteContext interface, which exposes Flink's table planner resources (such as the catalog reader, type factory, and function catalog) to external dialect implementations, and adds the PlannerExternalQueryOperation class to wrap Calcite RelNode trees and resolved schemas for use in the planning phase.

flink-table/flink-table-calcite-bridge · high confidence

Introduce test-filesystem connector and catalog for materialized table testing

Adds a new test-filesystem connector and catalog implementation to the test-utils module, enabling users to create and manage tables backed by the local filesystem for testing materialized table features. The new \TestFileSystemCatalog\ and \TestFileSystemTableFactory\ support both streaming and batch read modes, handle partition fields, and allow creating generic tables, providing a concrete file-based backend for validating materialized table logic in tests.

flink-test-utils-parent/flink-table-filesystem-test-utils · high confidence

Introduces generic AsyncSinkBase framework for async destinations

Adds a new extensible base class, AsyncSinkBase, and its supporting writer components (AsyncSinkWriter, AsyncSinkBaseBuilder, ElementConverter, BatchCreator) to the connector-base module. This framework allows connector developers to quickly implement sinks for destinations that support asynchronous ingestion, handling batching, buffering, and at-least-once semantics via checkpointing. It also includes a configurable exception classification system (FatalExceptionClassifier) to distinguish between retryable and fatal errors, and a state serializer (AsyncSinkWriterStateSerializer) to persist buffered requests across restarts.

flink-connectors/flink-connector-base/src/main · high confidence

Introduces the SQL Gateway API interface and supporting types

This change adds the core API contract for the SQL Gateway, introducing the SqlGatewayService interface which defines methods for session management (open, close, configure), operation handling (submit, cancel, close, get info), and statement execution (execute, fetch results). It also includes the necessary supporting types such as SessionHandle, OperationHandle, OperationStatus, ResultSet, and various info classes (FunctionInfo, TableInfo, GatewayInfo), along with configuration options for session timeouts, plan caching, and worker threads, and the endpoint factory framework for discovering and creating gateway endpoints.

flink-table/flink-sql-gateway-api · high confidence

Introduction of API stability annotations

The flink-annotations module now includes new Java annotations to help developers understand the stability and intended usage of Flink's public interfaces. Specifically, @Public marks stable interfaces that will not change within a major release, @PublicEvolving marks public APIs with stable behavior but evolving signatures, @Experimental marks classes that are not yet stable and may be changed or removed, @Internal marks stable but internal developer APIs, and @VisibleForTesting marks members that are only accessible for testing purposes.

flink-annotations/src/main/java/org/apache/flink/annotation · high confidence

Introduction of FlinkVersion enum for API versioning and migration

A new \FlinkVersion\ enum has been added to the \flink-annotations\ module to serve as a global registry for Flink versions. This enum includes versions from 1.3 through 2.4 and provides utilities for version arithmetic, range queries, and code-to-version mapping. It is designed to support API versioning, SQL/Table API upgrades, and state migration tests by ensuring consistent version identification across the platform.

flink-annotations/src/main/java/org/apache/flink · high confidence

This change introduces the new flink-rpc-core module, which extracts and centralizes the core abstractions for Flink's Remote Procedure Call (RPC) system. For users and developers, this provides a standardized, implementation-agnostic foundation for RPC endpoints and services. The module defines key interfaces such as RpcService for managing RPC connections and servers, RpcEndpoint as the base class for distributed components, and RpcGateway for accessing remote procedures. It also introduces support for fenced RPCs via FencedRpcEndpoint and FencedRpcGateway to ensure message ordering and safety, utilities for classloading isolation (ClassLoadingUtils) to prevent plugin classloader leaks, and main-thread execution guarantees (MainThreadExecutable, ComponentMainThreadExecutor) to maintain single-threaded endpoint execution semantics. Additionally, it includes infrastructure for resource cleanup (CleanupOnCloseRpcSystem) and fatal error handling (FatalErrorHandler).

flink-rpc/flink-rpc-core · high confidence

Major SQL parser overhaul with extensive new syntax support

The SQL parser has been significantly restructured and expanded to support a wide range of new Flink SQL capabilities. This includes comprehensive Data Definition Language (DDL) support for catalogs, databases, connections, functions, models, and materialized tables (including CREATE, ALTER, DROP, and SHOW variants). New Data Manipulation Language (DML) features include support for INSERT OVERWRITE, TRUNCATE, and statement sets. The parser now handles advanced query features like EXPLAIN PLAN ADVICE, MODEL syntax for machine learning, MATCH\_RECOGNIZE, and TRY\_CAST. Additionally, new data types such as VARIANT, BITMAP, and RAW are supported, along with enhanced SHOW commands for listing catalogs, databases, tables, views, functions, jars, modules, and partitions. The underlying parser implementation has been updated to use a new code generation structure (FMPP/FreeMarker) and imports numerous new Flink-specific SQL node classes to facilitate these changes.

flink-table/flink-sql-parser · high confidence

The \flink-architecture-tests-base\ module now includes a new set of common utilities to simplify writing architecture rules. This adds \Conditions\ for checking method leaf types (return, argument, and exception), \JavaFieldPredicates\ for inspecting field modifiers and assignability, and \Predicates\ for common field checks (e.g., public static final). It also introduces \GivenJavaClasses\ to restrict rules to Java classes (avoiding Scala issues), \ImportOptions\ to filter Maven main classes and exclude Scala/shaded code, and \SourcePredicates\ to identify Java sources. These helpers allow test authors to write more precise and maintainable architectural constraints.

flink-architecture-tests/flink-architecture-tests-base/src/main · high confidence

New Avro format factories and bulk I/O components for the Table API

The flink-avro module introduces new factory classes (AvroFormatFactory and AvroFileFormatFactory) that implement the modern Table API interfaces (DeserializationFormatFactory, SerializationFormatFactory, BulkReaderFormatFactory, and BulkWriterFormatFactory), enabling Avro support in the Table & SQL API. This change adds new bulk I/O components (AbstractAvroBulkFormat, AvroBulkWriter, AvroInputFormat, AvroOutputFormat) and configuration options (AvroFormatOptions) to handle Avro serialization and deserialization, including support for binary and JSON encoding, compression codecs, and legacy timestamp mapping.

flink-formats/flink-avro · high confidence

The \tools/azure-pipelines\ directory now contains a complete set of Azure Pipelines YAML templates and helper scripts to orchestrate Flink's continuous integration. This includes definitions for PR-triggered CI builds, nightly cron builds (including binary releases, Maven snapshot deployments, and Python wheel builds), and end-to-end test execution. The configuration introduces a modular structure with reusable templates for compilation, testing, and QA checks, utilizing a shared Docker container image and caching strategies for Maven repositories and E2E artifacts to improve build stability and speed.

tools · high confidence

New CI tools for license, dependency, and Scala suffix compliance

The flink-ci-tools module now includes new utilities to enforce build compliance: JarFileChecker validates that produced JARs contain valid NOTICE and LICENSE files and flags incompatible license text; ShadeOptionalChecker ensures dependencies bundled via the shade plugin are marked optional in POMs to maintain correct transitivity; ScalaSuffixChecker detects modules that incorrectly depend on Scala libraries; and parsers for Maven dependency, deploy, shade, and NOTICE outputs support these checks. These tools help catch licensing and dependency-management issues early in CI.

tools/ci/flink-ci-tools · high confidence

New DataGen and Blackhole connectors in the Table API Java Bridge

The \flink-table-api-java-bridge\ module now includes the \datagen\ and \blackhole\ connectors as first-class Table API components. The \datagen\ connector allows users to create tables that generate random or sequential data, supporting options for row count, emission rate, and per-field configuration (such as min/max bounds, null rates, and variable lengths). The \blackhole\ connector provides a sink that discards all input records, useful for performance testing or as a UDF output target. Both connectors are implemented using the modern \DynamicTableSource\/\DynamicTableSink\ factory interfaces.

flink-table/flink-table-api-java-bridge · high confidence

New DataStream API Java walkthrough for fraud detection

A new Maven archetype has been added to generate a Java-based DataStream API walkthrough project. This scaffold provides a complete example application that reads an infinite stream of credit card transactions, applies a keyed process function to detect fraudulent patterns, and outputs alerts. The generated project includes the necessary source files (FraudDetectionJob, FraudDetector) and is pre-configured with a Log4j2 logging setup to help users quickly start building and experimenting with Flink DataStream applications.

flink-walkthroughs/flink-walkthrough-datastream-java · high confidence

New DataStream API V2 and Async IO examples added

The streaming examples module now includes new demonstration programs for the experimental DataStream API V2, covering WordCount, event-time processing, joins, watermarks, and windowing, alongside new examples for asynchronous I/O operations and data generation sources.

flink-examples/flink-examples-streaming · high confidence

New DataStream Allround End-to-End Test Job

A new general-purpose end-to-end test job has been added to validate Flink's DataStream API operators and primitives. The \DataStreamAllroundTestProgram\ and its factory \DataStreamAllroundTestJobFactory\ construct a comprehensive pipeline that exercises keyed state (using both Kryo and Avro serializers), operator state, sliding and tumbling event-time windows, and configurable failure simulation to verify exactly-once and at-least-once processing semantics. The test suite also includes support for externalized checkpoints, restart strategies, and various state backend configurations.

flink-end-to-end-tests/flink-datastream-allround-test · high confidence

New Hadoop compatibility wrappers for mapred APIs

The flink-hadoop-compatibility module now includes new utility classes and wrappers to bridge Apache Hadoop mapred APIs with Flink. HadoopInputs provides static methods to create Flink InputFormat wrappers for Hadoop FileInputFormat and SequenceFileInputFormat. HadoopUtils offers a helper to parse command-line arguments using Hadoop's GenericOptionsParser into Flink's ParameterTool. Additionally, HadoopMapFunction wraps Hadoop mappers into Flink FlatMapFunctions, and HadoopReducerWrappedFunction wraps Hadoop reducers into Flink window functions (both keyed and non-keyed). Supporting classes HadoopOutputCollector and HadoopTupleUnwrappingIterator facilitate the data exchange between Hadoop and Flink runtime components.

flink-connectors/flink-hadoop-compatibility/src/main/java/org/apache/flink/hadoopcompatibility · high confidence

New JUnit 5 test utilities and assertion helpers

The flink-test-utils-junit module now provides a suite of new utilities to support JUnit 5-based testing. This includes JUnit 5 extension wrappers (AllCallbackWrapper, EachCallbackWrapper, CustomExtension) to simplify lifecycle management, synchronization aids (BlockerSync, MultiShotLatch, OneShotLatch), and a CheckedThread class for safer asynchronous test execution. Additionally, new AssertJ-based assertions (FlinkAssertions, FlinkCompletableFutureAssert) allow for more expressive checks on exception chains and CompletableFuture states without relying on timeouts, alongside a ManuallyTriggeredScheduledExecutorService for deterministic scheduling tests.

flink-test-utils-parent/flink-test-utils-junit · high confidence

New OpenTelemetry reporters for metrics, traces, and events

This release introduces three new OpenTelemetry reporters for Apache Flink: OpenTelemetryMetricReporter, OpenTelemetryTraceReporter, and OpenTelemetryEventReporter. These components allow Flink to export standard metrics, distributed traces, and internal system events to any OpenTelemetry-compatible backend. The reporters support both gRPC and HTTP protocols, configurable endpoint and timeout settings, and optional gzip compression. Additionally, the metric reporter includes configurable attribute value length limits to control data size and a collision-tracking mechanism to detect when truncation causes distinct metrics to merge, helping to maintain data integrity in downstream observability systems.

flink-metrics/flink-metrics-otel · high confidence

New ProcessTableFunction test harness with state and timer support

The flink-table-test-utils module now includes a ProcessTableFunctionTestHarness that allows testing ProcessTableFunctions with table and scalar arguments, lifecycle management, and output collection. The harness supports state management (ListView, MapView, Row, and structured types) and timer management (registration, clearing, watermark tracking, and firing). It also provides a fluent builder API for configuring and testing ProcessTableFunctions.

flink-table/flink-table-test-utils · high confidence

The end-to-end test suite now includes dedicated test harnesses for PyFlink DataStream applications and TPC-DS benchmarks. A new \flink-python-test\ module provides Python scripts (e.g., \data\_stream\_job.py\) that exercise modern PyFlink APIs such as \FileSource\, \FileSink\, and \KeyedProcessFunction\ with timers, alongside a word-count example using Python UDFs. Additionally, the \flink-tpcds-test\ module introduces the TPC-DS benchmark tooling, including data generation scripts, 103 standard SQL queries, and expected answer sets, enabling validation of Flink's SQL engine against a complex analytical workload.

flink-end-to-end-tests · high confidence

New REST API for managing uploaded JAR files

The Flink Web Dashboard now exposes a dedicated REST API for managing uploaded JAR files, allowing users to list, run, plan, and delete JARs directly via HTTP requests. This change introduces new REST handlers (JarListHandler, JarRunHandler, JarPlanHandler, JarDeleteHandler) and their corresponding message headers, which are registered through the WebSubmissionExtension. The API endpoints support operations such as retrieving the list of uploaded JARs, triggering job execution with optional parameters like entry class and savepoint restore settings, generating execution plans, and removing JAR files from the server's storage directory.

flink-runtime-web · high confidence

New and updated Table API examples

The examples-table module now includes several new Java examples demonstrating modern Table and SQL API usage: a basic GettingStartedExample for batch processing, StreamSQLExample for streaming SQL on DataStreams, StreamWindowSQLExample for windowed aggregations, TemporalJoinSQLExample for time-based joins, UpdatingTopCityExample for top-N ranking on changelogs, and WordCountSQLExample as a minimal batch SQL job. Additionally, a ChangelogSocketExample with custom connector components (SocketDynamicTableFactory, ChangelogCsvFormat, ChangelogCsvDeserializer) illustrates how to implement and use custom DynamicTableSources and DecodingFormats with changelog semantics.

flink-examples/flink-examples-table · high confidence

New automated documentation generators for REST APIs and configuration options

The flink-docs module now includes new Java-based generators that automatically produce reference documentation directly from the Flink source code. The ConfigOptionsDocGenerator scans configured modules and packages to generate HTML tables for all ConfigOptions, handling grouping, sections, and exclusions. The REST API documentation is now generated via RestAPIDocGenerator (producing HTML files) and OpenApiSpecGenerator (producing OpenAPI YAML specs) for the Runtime and SQL Gateway endpoints, ensuring the docs stay in sync with endpoint handlers, message headers, and request/response schemas.

flink-docs · high confidence

New compression support for FileSink via CompressWriterFactory

The flink-compress module now provides a new CompressWriterFactory that enables users to write compressed bulk files using the modern FileSink API. This factory wraps Hadoop CompressionCodecs (such as Gzip, Bzip2, and Deflate) to compress data during the write process, supporting both standard and custom codecs via configuration. It includes a DefaultExtractor for simple string-based records and is designed to integrate seamlessly with the FileSink's BulkFormat capability, allowing users to produce compressed output files without relying on the deprecated StreamingFileSink.

flink-formats/flink-compress · high confidence

New core API interfaces for DataStream V2 state and execution modes

The flink-core-api module now exposes the foundational interfaces for the DataStream V2 API, including RuntimeExecutionMode (STREAMING, BATCH, AUTOMATIC) to control pipeline semantics, and a comprehensive set of state interfaces (State, ListState, MapState, BroadcastState, ReducingState, AggregatingState) along with their corresponding declarations (StateDeclaration, ListStateDeclaration, etc.) that allow users to define and manage partitioned state with explicit type descriptors and redistribution strategies. Additionally, base functional interfaces such as Function, ReduceFunction, and AggregateFunction have been moved into this module to support the new API's user-defined transformations.

flink-core-api · high confidence

New core metrics interfaces and event reporting infrastructure

The metrics core module now exposes a comprehensive set of new public interfaces and classes that expand the observability capabilities of Flink. This includes the foundational metric types (Counter, Gauge, Histogram, Meter) and their implementations (SimpleCounter, MeterView, SlidingWindowHistogram), alongside a new Event reporting system (Event, EventBuilder, EventReporter) that allows external backends to consume structured events like job status changes and checkpoint completions. Additionally, the module introduces MetricConfig for robust configuration handling, CharacterFilter for metric name sanitization, and LogicalScopeProvider for better metric grouping context, all marked as Experimental or Public to support future plugin and reporter development.

flink-metrics/flink-metrics-core · high confidence

New state migration test framework and snapshot generator

The flink-migration-test-utils module now includes a new framework for generating reference data for state migration tests. This introduces the MigrationTest interface with @SnapshotsGenerator and @ParameterizedSnapshotsGenerator annotations, a MigrationTestsSnapshotGenerator CLI tool to scan and generate snapshots for specified Flink versions, and a JUnit5TestEnvironment to drive JUnit 5 lifecycle methods (extensions, @BeforeEach/@AfterEach, @TempDir) during snapshot generation outside the standard JUnit engine. The module also includes PublishedVersionUtils to read the most recently published Flink version (currently v2.3) from a resource file.

flink-test-utils-parent/flink-migration-test-utils · high confidence

The flink-end-to-end-tests-common module introduces a suite of new utility classes to support end-to-end test execution. This includes AutoClosableProcess for managing external process lifecycles with configurable timeouts, retries, and I/O handling, and AutoClosablePath for automatic cleanup of temporary files. A new DownloadCache system (with implementations like PersistingDownloadCache and TravisDownloadCache) allows tests to cache downloaded artifacts with configurable time-to-live policies, managed via a JUnit 5 extension. Additionally, CommandLineWrapper provides fluent builders for common shell commands (wget, sed, tar), and TestUtils offers helpers for directory copying and CSV reading.

flink-end-to-end-tests/flink-end-to-end-tests-common · high confidence

New test utilities for Source V2 implementations

The flink-test-utils-connector module now includes a suite of classes to simplify writing tests for Flink Source V2 connectors. This adds abstract base classes (AbstractTestSource, AbstractTestSourceBase) that provide default implementations for Source, SourceReader, and SplitEnumerator methods, allowing test authors to override only the specific behaviors they need, such as data emission in pollNext. It also includes concrete test utilities like TestSplit, TestSplitEnumerator, SingleSplitEnumerator, TestSourceReader, and TestReaderOutput to handle split management, enumerator logic, and output verification without implementing the full interfaces from scratch.

flink-test-utils-parent/flink-test-utils-connector · high confidence

The flink-test-utils module now includes several new components to aid in testing Flink applications. This includes an UpsertTestSink connector (with a corresponding DynamicTableSink and factory) that writes upsert-style data to files for verification, a MetricListener and MetricAssertions class for AssertJ-based metric testing, a FiniteTestSource for controlled stream emission, and PackagingTestUtils for validating JAR contents.

flink-test-utils-parent/flink-test-utils · high confidence

New typed builders for DefaultRollingPolicy and introduction of OnCheckpointRollingPolicy

The file sink common module now provides a new typed builder API for configuring the DefaultRollingPolicy, replacing the previous untyped create() method with a fluent builder that accepts MemorySize and Duration types for part size, rollover interval, and inactivity interval. Additionally, a new OnCheckpointRollingPolicy class has been added, which implements a rolling policy that triggers file rolls exclusively on checkpoints, ignoring event count, processing time, and inactivity triggers.

flink-connectors/flink-file-sink-common/src/main/java/org/apache/flink/streaming/api/functions/sink/filesystem/rollingpolicies · high confidence

New utility class for parsing and applying JAR request parameters in the REST API

Added JarHandlerUtils to centralize the extraction of JAR handler parameters (such as entry class, program arguments, parallelism, and job ID) from REST requests and their application to the Flink configuration. This utility supports both standard JAR run/plan requests and the new application-level requests (JarRunApplicationRequestBody), ensuring that pipeline classpaths and configuration settings are correctly propagated during job deployment.

flink-runtime-web/src/main/java/org/apache/flink/runtime/webmonitor/handlers/utils · high confidence

Prometheus reporters now support HTTP Basic Authentication for PushGateway

The Prometheus PushGateway reporter now supports HTTP Basic Authentication when pushing metrics to a secured PushGateway. Users can configure the \username\ and \password\ options to enable authentication; if only one is provided, a warning is logged and authentication remains disabled. The implementation uses a JDK-based encoder to avoid the JAXB dependency previously required by the Prometheus client library. Additionally, the reporter now supports a \hostUrl\ configuration option to specify the PushGateway server address, and includes a \randomJobNameSuffix\ option to append a unique identifier to the job name for better metric isolation.

flink-metrics/flink-metrics-prometheus · high confidence

PyFlink has added support for Python 3.12, allowing users to run their Python DataStream and Table API jobs using this newer Python version.

flink-python · high confidence

Raw format now supports line delimiters and VARIANT type

The 'raw' format in the Table Runtime now supports an optional line-delimiter configuration option, allowing deserialization to split incoming byte messages into multiple rows and serialization to append delimiters after each row. Additionally, the format's supported column types have been extended to include VARIANT, enabling direct reading and writing of variant data in raw format sources and sinks.

flink-table/flink-table-runtime · high confidence

State Processor API now supports reading and writing evicting window state

The State Processor API in flink-state-processing-api has been extended to support reading and writing state from evicting window operators. Users can now use the new EvictingWindowSavepointReader to read keyed state generated by ReduceFunction, AggregateFunction, or ProcessFunction-based evicting windows, and the StateBootstrapTransformation API to bootstrap new evicting window state into savepoints via the SavepointWriter.

flink-libraries/flink-state-processing-api · high confidence

Support for Hadoop path-based part-file writing in streaming sinks

This change introduces a new path-based bulk writer implementation for Hadoop filesystems, allowing streaming sinks to write directly to specified Hadoop paths rather than relying on output-stream-based writers. It adds the \HadoopPathBasedBulkWriter\ interface and the \HadoopPathBasedPartFileWriter\ class to handle writing and committing files via a \HadoopFileCommitter\ (specifically \HadoopRenameFileCommitter\, which uses UUIDs for temporary files to avoid conflicts). A new \HadoopPathBasedBulkFormatBuilder\ is provided to configure and create buckets using this path-based approach, and the \PendingFileRecoverable\ serializer is updated to version 2 to support the new file size tracking.

flink-formats/flink-hadoop-bulk · high confidence

Unified FileSink with built-in file compaction

The FileSystem connector now provides a new unified FileSink that supports exactly-once semantics for both batch and streaming workloads. This sink introduces built-in file compaction capabilities, allowing users to merge small output files into larger ones to improve downstream read performance. The implementation includes a configurable compaction strategy (triggered by size thresholds or checkpoint counts), a coordinator to manage compaction requests, and operators to handle the actual merging of files, all integrated directly into the sink's commit process.

flink-connectors/flink-connector-files/src/main · high confidence

Architecture

File sink common module refactored with new internal writer interfaces and bucket assigners

The file sink common module has been restructured to introduce a new internal architecture for file writing and bucketing. This includes new interfaces such as InProgressFileWriter, CompactingFileWriter, and BucketWriter to manage part file lifecycle and recovery, alongside implementations like RowWisePartWriter and BulkPartWriter for different encoding strategies. Additionally, new bucket assigners (BasePathBucketAssigner, DateTimeBucketAssigner) and configuration classes (OutputFileConfig, RollingPolicy) have been added to control how data is partitioned into files. These changes are internal to the file sink implementation and do not expose new public APIs, but they provide the foundation for more robust and flexible file sinking behavior.

flink-connectors/flink-file-sink-common/src/main/java/org/apache/flink/streaming/api/functions/sink/filesystem · high confidence

This change introduces the new flink-table-planner-loader module, which acts as a delegation layer for the Table API's Executor and Planner factories. By loading the planner implementation dynamically via the PlannerModule and registering DelegateExecutorFactory and DelegatePlannerFactory in the service provider configuration, this module ensures that the heavy flink-table-planner jar is not bundled in the client classpath by default. This architectural shift prevents the planner jar from leaking into temporary directories (e.g., /tmp) during SQL client operations, thereby reducing client footprint and avoiding classpath pollution.

flink-table/flink-table-planner-loader · high confidence

New type-utils module centralizes data structure conversion

The new flink-table-type-utils module introduces a unified DataStructureConverter registry that standardizes how internal Flink data structures (like RowData and ArrayData) are converted to and from external Java types (such as primitive arrays, java.time classes, and collections). This refactoring consolidates previously scattered conversion logic into a single, reusable component, ensuring consistent handling of type conversions across the Table API and enabling more reliable data exchange at API boundaries.

flink-table/flink-table-type-utils · high confidence

RocksDB state backend classes relocated to a new package

The core RocksDB state backend implementation classes, including EmbeddedRocksDBStateBackend, AbstractRocksDBState, and various state-specific implementations, have been moved from the org.apache.flink.contrib.streaming.state package to the new org.apache.flink.state.rocksdb package. To ensure backward compatibility, deprecated shim classes and interfaces remain in the original location and delegate to the new implementations, allowing existing code to continue functioning while encouraging migration to the new package structure.

flink-state-backends/flink-statebackend-rocksdb · high confidence

SQL Client refactored into a unified gateway and embedded architecture

The SQL Client has been restructured to support both embedded and gateway execution modes through a new entry point (SqlClient) and a decoupled CLI layer (CliClient). This change introduces a new exception class (SqlClientException) and a dedicated changelog result view (CliChangelogResultView) for streaming results, while consolidating command-line option parsing (CliOptions, CliOptionsParser) and input handling (CliInputView) into a cleaner, reusable structure.

flink-table/flink-sql-client · high confidence

Behavioural changes

The flink-sql-orc module now includes the required META-INF/NOTICE file, which lists bundled dependencies such as Apache ORC, Hive Storage API, and Google Protocol Buffers, along with the specific BSD license file (LICENSE.protobuf) for the protobuf dependency. This ensures proper attribution and legal compliance for third-party components included in the distribution.

flink-formats/flink-sql-orc · high confidence

Added license files for anchorjs, chroma, cloudpickle, font-awesome, and py4j

The licenses directory now includes explicit license texts for anchorjs, chroma, cloudpickle, font-awesome, and py4j. This ensures that the legal terms for these dependencies are clearly documented and distributed alongside the software.

licenses · high confidence

Consolidated S3 filesystem implementation with delegation token and s5cmd support

The flink-s3-fs-base module now provides a unified, consolidated implementation of the S3 filesystem, replacing scattered code with a single, coherent set of classes. This change introduces S3 delegation token support via AbstractS3DelegationTokenProvider and AbstractS3DelegationTokenReceiver, enabling dynamic session credentials through DynamicTemporaryAWSCredentialsProvider. It also integrates s5cmd for faster file transfers during operations like RocksDB incremental state recovery, configurable via new options such as s3.s5cmd.path and s3.s5cmd.adjust-part-size. The filesystem factory (AbstractS3FileSystemFactory) and main filesystem class (FlinkS3FileSystem) are streamlined, and common utilities like S3AccessHelper centralize S3 access logic, improving maintainability and performance.

flink-filesystems/flink-s3-fs-base · high confidence

Default batch mode uses blocking global stream exchange

In batch mode, the streaming engine now defaults to the ALL\_EDGES\_BLOCKING shuffle mode (GlobalStreamExchangeMode). This ensures that all upstream operators complete and flush their data before downstream operators begin processing, providing the deterministic, pull-based execution semantics required for batch workloads.

flink-streaming-java · high confidence

Deprecate complex Java objects in restart strategy configuration

The configuration options for restart strategies that previously accepted complex Java objects (such as specific strategy instances) are now deprecated. Users should instead use the standard configuration keys and values (e.g., string-based strategy names and numeric parameters) to define restart behavior, ensuring configuration is portable and easier to manage across cluster deployments.

flink-core · high confidence

Deprecated Scala Table API bridge for DataStream conversions

The Scala bridge module now provides the \StreamTableEnvironment\, \DataStreamConversions\, and \TableConversions\ classes, which enable implicit conversions between Scala \DataStream\s and \Table\s (including changelog streams) and support attaching statement sets to the execution environment. These APIs are marked as deprecated since version 1.18.0 under FLIP-265, signaling that users should migrate to the Java version of the DataStream and/or Table API in future major versions.

flink-table/flink-table-api-scala-bridge · high confidence

Deprecated legacy Table API methods moved to legacy package

The \flink-table-api-java\ module has relocated deprecated user-visible classes and methods to a dedicated \legacy\ package. This change isolates older, unsupported APIs from the main public interface, signaling that they are scheduled for removal in a future version and encouraging users to migrate to the current Table API surface.

flink-table/flink-table-api-java · high confidence

Enforced architectural rules for test code quality and JUnit 5 migration

The test code architecture module now enforces several new rules to standardize test structure and accelerate the JUnit 5 migration. Tests inheriting from AbstractTestBase must end with 'ITCase', and all integration tests must use the JUnit 5 MiniClusterExtension or the legacy MiniClusterWithClientResource. To prevent silent test failures, static nested test classes must be non-static inner classes annotated with @Nested, and all executable test classes must be named with 'Test', 'Tests', or 'ITCase' suffixes to satisfy Surefire patterns. Additionally, modules that have completed the JUnit 5 migration are now forbidden from adding new JUnit 4 dependencies.

flink-architecture-tests/flink-architecture-tests-test/src/main · high confidence

Hadoop S3 connector now supports the s3a:// scheme and standard YAML configuration

The Hadoop-based S3 filesystem plugin now registers a factory for the s3a:// scheme (in addition to the existing s3:// support), allowing users to access S3 buckets using the Hadoop-native s3a URI prefix. The S3FileSystemFactory also introduces configuration key aliases with a .value suffix (e.g., fs.s3a.endpoint.value) to resolve conflicts in Flink v2's standard YAML format, where a key cannot be both a scalar value and a parent of other keys. Additionally, the connector consolidates S3 access logic into a new HadoopS3AccessHelper class to streamline multipart uploads and object operations.

flink-filesystems/flink-s3-fs-hadoop · high confidence

Introduce KubernetesResourceManagerDriver and new Kubernetes deployment components

The flink-kubernetes module now includes the KubernetesResourceManagerDriver, which manages TaskManager pod lifecycle and resource requests via the Kubernetes API, replacing the previous resource manager implementation. This change is accompanied by the introduction of KubernetesClusterClientFactory and KubernetesClusterDescriptor to handle cluster deployment and retrieval, a new KubernetesArtifactUploader interface with a Default implementation for uploading local artifacts to a remote target, and a KubernetesSessionCli for managing session clusters. These components collectively enable the new active resource manager architecture for Kubernetes deployments.

flink-kubernetes · high confidence

Introduce common base classes for Hadoop input and output formats

Added HadoopInputFormatCommonBase and HadoopOutputFormatCommonBase to centralize the serialization and deserialization of Hadoop security credentials (Credentials) for both mapred and mapreduce formats. This refactoring ensures consistent handling of credential objects across Hadoop input and output format implementations within the Flink Hadoop compatibility layer.

flink-connectors/flink-hadoop-compatibility/src/main/java/org/apache/flink/api/java/hadoop/common · high confidence

JMX metrics reporter now uses a factory-based SPI registration

The JMX metrics reporter in flink-metrics-jmx has been refactored to use a service-provider interface (SPI) approach for instantiation. A new JMXReporterFactory class and its corresponding META-INF/services registration file replace the previous direct instantiation model, allowing the reporter to be configured and loaded dynamically via the standard Java SPI mechanism. This change also introduces a new JMXReporterFactoryTest to verify the factory's behavior and updates the JMXJobManagerMetricTest to use the new factory-based configuration, ensuring that metrics are correctly exposed via JMX when the reporter is loaded through the service loader.

flink-metrics/flink-metrics-jmx · high confidence

JSON format now uses a high-performance JsonParser-based deserialization schema by default

The JSON format in the Table/SQL API now defaults to a new \JsonParser\-based deserialization path (\JsonParserRowDataDeserializationSchema\) instead of the previous Jackson \JsonNode\-based approach. This change, controlled by the \decode.json-parser.enabled\ option (defaulting to \true\), significantly improves parsing performance and enables projection push-down for nested fields. The new implementation also introduces support for the \VARIANT\ type, allows parsing of unescaped control characters in strings, and provides better error logging for parse failures. Users relying on the old behavior can disable the new parser by setting \decode.json-parser.enabled\ to \false\.

flink-formats/flink-json · high confidence

Modernized Java quickstart archetype with Log4j2 and updated build structure

The Java quickstart archetype has been modernized to generate projects using Log4j2 for logging (replacing the previous log4j.properties configuration) and includes an updated archetype descriptor and test resources to ensure the generated project builds correctly.

flink-quickstart/flink-quickstart-java · high confidence

New HadoopUtils helper for loading Hadoop configuration

A new HadoopUtils utility class has been added to the Hadoop compatibility module to centralize logic for loading and merging Hadoop configurations. This class provides methods to retrieve possible Hadoop configuration paths by checking environment variables (HADOOP\_CONF\_DIR, HADOOP\_HOME) and Flink configuration, and to merge those settings into a JobConf instance. This supports the HadoopInputFormatBase in accessing necessary HDFS configuration details.

flink-connectors/flink-hadoop-compatibility/src/main/java/org/apache/flink/api/java/hadoop/mapred/utils · high confidence

New architecture rules enforce API visibility, connector dependencies, and checkpointing configuration access

This change introduces new ArchUnit rules in the production architecture test module to enforce stricter coding standards. API visibility rules now require all public classes in API packages to carry a visibility annotation (e.g., @Public, @Internal) and ensure that methods annotated with @Public or @PublicEvolving only use types with matching visibility annotations. Connector production code is now restricted to depending only on public Flink API classes and internal utilities, preventing accidental reliance on internal APIs. Additionally, a new rule prevents direct access to specific checkpointing configuration options (like ENABLE\_UNALIGNED) via Configuration.get() or getOptional(), mandating the use of dedicated helper methods instead. Finally, production code is prohibited from calling methods annotated with @VisibleForTesting.

flink-architecture-tests/flink-architecture-tests-production/src/main/java/org/apache/flink/architecture/rules · high confidence

The flink-table module has been reorganized into a new structure with dedicated modules for common components, APIs, runtime, parser/planner, SQL client, and testing utilities. This change includes new README and AGENTS documentation files for the table planner and runtime modules, providing detailed guidance on module responsibilities, directory structures, and development patterns. The SQL client and SQL gateway scripts have been updated to properly handle classpath management and JVM options.

flink-table · high confidence

New serializer snapshots for Scala types to support state migration

The Scala Table API now includes dedicated TypeSerializerSnapshot classes for Scala-specific types such as Option, Try, Either, Traversable, and Case Classes. These snapshots ensure that the internal configuration of these serializers is correctly preserved and resolved during state migration, allowing applications using these Scala types to upgrade Flink versions without serialization compatibility issues.

flink-table/flink-table-api-scala · high confidence

ORC format reader refactored to use the new FileSource API with vectorized columnar reading

The ORC input format has been rewritten to integrate with the new FileSource connector API, replacing the legacy input format with a new set of classes including AbstractOrcFileInputFormat, OrcColumnarRowInputFormat, and OrcSplitReader. This change enables vectorized columnar reading for improved performance and introduces support for filter push-down via the new OrcFilters utility class. The implementation also adds a shim layer (OrcShim) to maintain compatibility with different Hive versions (2.0, 2.1, 2.3, and 3.x) and ensures proper handling of timestamp vectors across these versions.

flink-formats/flink-orc · high confidence

Optimize Hadoop mapred input split serialization

The Hadoop mapred compatibility layer now optimizes serialization by only serializing the JobConf when the input split actually requires it (i.e., when it implements Configurable or JobConfigurable). This change introduces new wrapper classes (HadoopDummyProgressable, HadoopDummyReporter, HadoopInputSplit) to support this behavior, reducing unnecessary data transfer and improving performance for splits that do not need configuration context.

flink-connectors/flink-hadoop-compatibility/src/main/java/org/apache/flink/api/java/hadoop/mapred/wrapper · high confidence

Queryable State Client API restructured with immutable state wrappers

The Queryable State Client module has been restructured to improve API clarity and safety. The client now returns read-only, immutable wrappers for all queried state types (ValueState, ListState, MapState, AggregatingState, and ReducingState), preventing accidental modifications to the retrieved data. Additionally, the module has been decoupled from the runtime module by introducing local copies of the VoidNamespace class and its serializer, eliminating a direct dependency on the runtime package for this client component.

flink-queryable-state/flink-queryable-state-client-java · high confidence

Queryable State runtime restructured with new proxy and server implementations

The queryable state runtime module has been restructured to introduce new internal components for handling state queries. This includes a new \KvStateClientProxyHandler\ and \KvStateClientProxyImpl\ that act as an internal proxy to receive requests, resolve state locations via a \KvStateLocationOracle\, and forward queries to the appropriate Task Manager. A new \KvStateServerHandler\ and \KvStateServerImpl\ handle the actual state retrieval from the \KvStateRegistry\ on the Task Manager side. The communication between the proxy and the server uses a new \KvStateInternalRequest\ message type. These changes are accompanied by updated integration tests (\HAQueryableStateFsBackendITCase\, \HAQueryableStateRocksDBBackendITCase\, etc.) that verify the functionality of the queryable state feature with both file system and RocksDB state backends.

flink-queryable-state/flink-queryable-state-runtime · high confidence

RPC layer migrated from Akka to Apache Pekko

The Flink RPC implementation in flink-rpc-akka has been rewritten in Java to use Apache Pekko instead of Akka. This migration introduces new core components such as PekkoRpcActor, PekkoInvocationHandler, and PekkoBasedEndpoint, alongside supporting utilities like ActorSystemScheduledExecutorAdapter and ScalaFutureUtils to bridge Scala/Java types. The change also brings enhanced security features, including a CustomSSLEngineProvider that supports comma-separated SSL protocols and configurable keystore/truststore types, as well as improved fault handling via a DeadLettersActor that proactively fails RPC requests for unreachable recipients.

flink-rpc/flink-rpc-akka · high confidence

Redesigned binary distribution assembly and startup configuration

The Flink binary distribution has been restructured to use new Maven assembly descriptors (bin.xml, opt.xml, plugins.xml) that explicitly package core components, optional modules (such as the SQL client, state processor API, and various filesystem connectors), and metrics plugins into their respective directories. The startup scripts have been rewritten to rely on a new BashJavaUtils utility for dynamic configuration and memory calculation, and the default configuration file has been changed from the legacy flink-conf.yaml to config.yaml, requiring users to migrate their settings to the new YAML format.

flink-dist · high confidence

Rework test user jar assembly for classloading tests

The test utilities for Flink clients have been restructured to better simulate user-jar classloading scenarios. New test artifacts, including \TestUserClassLoaderJob\, \TestUserClassLoaderJobLib\, and \TestUserClassLoaderAdditionalArtifact\, are now provided to verify that jobs can correctly load classes from additional user artifacts that are not on the system classpath, ensuring robustness in distributed execution environments.

flink-test-utils-parent/flink-clients-test-utils · high confidence

Separate Scala distribution with explicit license and dependency notices

The flink-dist-scala module now includes a dedicated NOTICE file and a bundled LICENSE.scala file. The NOTICE file explicitly lists the Scala dependencies included in this distribution (scala-compiler, scala-library, scala-reflect, and scala-xml\_2.12) and their versions (2.12.21 and 2.3.0), while the LICENSE.scala file provides the full BSD license text required for redistribution. This ensures that users of the Scala-specific distribution package have the necessary legal attribution and license information for the included third-party Scala libraries.

flink-dist-scala · high confidence

SinkUpsertMaterializer no longer inserted for retract sinks

The table planner now correctly avoids inserting the SinkUpsertMaterializer for retract sinks, ensuring that retract-based sinks receive the expected changelog stream without unnecessary upsert materialization.

flink-table/flink-table-planner · high confidence

Support partial deletes when converting to external data structures

The table module now supports partial deletes when converting to external data structures, allowing users to emit only the changed fields in delete operations rather than full rows. This improves efficiency for downstream systems that can handle partial updates.

flink-table/flink-table-common · high confidence

Updated ArchUnit violation baselines for internal API exposure

The ArchUnit violation baselines in the architecture-tests-production module have been updated to reflect new architectural constraints. The diff shows the addition of two new violation files (identified by hashes 18509c9e and 5b9eed8a) that record violations where internal Flink classes (e.g., ExecutionConfig, DistributedCache, WatermarkStrategy, FileSystem, StateBackend) are exposed in public method signatures or arguments. These violations indicate that certain internal types are leaking into public APIs, which the ArchUnit rules now flag as violations of the requirement that public APIs must not depend on internal or shaded packages unless annotated as @Public, @PublicEvolving, or @Deprecated.

flink-architecture-tests/flink-architecture-tests-production/archunit-violations · high confidence

Updated JVM configuration for Java 17 compatibility and build stability

A new .mvn/jvm.config file has been introduced to configure the JVM arguments used during the build process. This includes enabling the -XX:+IgnoreUnrecognizedVMOptions flag to prevent build failures from unrecognized options, and adding specific --add-exports directives for java.security.jgss and jdk.compiler modules to ensure compatibility with Java 17's module system.

.mvn · high confidence

Updated Java Table API walkthrough to use modern connectors and UDFs

The Java Table API walkthrough archetype has been refreshed to demonstrate current best practices. The example application now uses the DataGen connector to generate streaming transaction data instead of the deprecated TableSource registration, and includes a custom ScalarFunction (MyFloor) to show how to implement user-defined functions for timestamp truncation. The project also migrates logging configuration from log4j to log4j2 and provides batch-mode unit tests to validate the aggregation logic.

flink-walkthroughs/flink-walkthrough-table-java · high confidence

Web dashboard restructured into standalone Angular components with new application navigation

The web dashboard has been refactored to use Angular standalone components, replacing the previous module-based architecture. This introduces a new application-centric navigation model with dedicated 'Running Applications' and 'Completed Applications' pages, alongside the existing job and cluster views. The update includes a new HTTP interceptor that handles authentication credentials, redirects, and server error notifications, along with a suite of reusable UI components such as status badges, compact action bars, and application lists. Additionally, TypeScript type definitions for the d3-flame-graph library have been added to support profiling visualizations.

flink-runtime-web/web-dashboard · high confidence

YARN resource allocation now matches containers by priority rather than resources

The YARN deployment module now uses a new adapter to map Flink TaskExecutor process specifications to YARN container priorities, changing the container matching strategy from resource-based to priority-based. This behavioral change ensures that requested and allocated containers are matched according to their priority levels, which aligns with the YARN FairScheduler's allocation logic and improves resource utilization in shared clusters.

flink-yarn · high confidence

Fixes

Add NOTICE and license files for ICU4j dependency

The flink-table-api-java-uber module now includes the required legal attribution files for its bundled dependencies. Specifically, a NOTICE file has been added to declare the Apache Software Foundation copyright (2014-2026) and list the bundled ICU4j library (com.ibm.icu:icu4j:67.1), accompanied by the corresponding ICU license text to ensure compliance with third-party licensing requirements.

flink-table/flink-table-api-java-uber · high confidence

Test coverage

1 commit adding/updating tests in flink-connectors/flink-hadoop-compatibility/src/test/resources/writeable-serializer-1.11; Add CLI test for periodic streaming job with checkpointing; Add Stream SQL end-to-end test program; Add end-to-end test for failure label enrichment; Add second dummy filesystem plugin for class-isolation testing; Add state evolution e2e test for Avro and complex types; Added ArchUnit configuration and migrated test logging to Log4j2; Added ArchUnit tests for test code architecture; Added E2E test for quickstart dependency packaging; Added JUnit 5 extension configuration for test utilities; Added JUnit 5 test infrastructure and migration tests for file sink recoverables; Added JUnit 5 tests for file enumerator implementations; Added JUnit 5 tests for the Flink client program execution layer; Added Java 17 record serialization and type-extraction tests; Added Java-based end-to-end test for batch SQL execution; Added YARN integration tests and test utilities; Added architecture tests for production code using ArchUnit; Added architecture tests for the files connector test code; Added end-to-end test for DistributedCache via BlobServer; Added end-to-end test for State TTL feature; Added end-to-end test for local recovery with sticky allocation; Added end-to-end test for metrics availability; Added end-to-end test for parent-child classloader resolution; Added end-to-end test for the JDBC driver; Added end-to-end tests for Python UDFs in batch and streaming modes; Added end-to-end tests for async scalar functions and joins with custom types; Added end-to-end tests for the SQL Gateway with Hive integration; Added integration and unit tests for the base source reader; Added integration test for REST client HTTPS connectivity; Added integration tests for Hadoop Input/Output formats using DataStream API; Added integration tests for the generic AsyncSinkBase implementation; Added log4j2 test logging configuration; Added packaging integration test and NOTICE file for flink-sql-avro; Added parameterizable testing mocks for the Split Reader API; Added test configuration for ArchUnit and logging; Added test fixtures for YARN client configuration validation; Added test infrastructure configuration files; Added test program for DataStream stateful job upgrade scenarios; Added test program for Netty shuffle direct memory control; Added test service provider configurations for JUnit 5 and ClusterClientFactory; Added test utilities for CLI client deployment mocking; Added test utilities for FileSink; Added test utilities for file source enumeration and file system mocking; Added test utility classes for client-side job submission testing; Added tests for ArtifactFetchManager and ArtifactUtils; Added tests for ClusterClientServiceLoader discovery behavior; Added tests for FatalExceptionClassifier; Added tests for FileSink compaction, serialization, and execution modes; Added tests for FileSink writer migration and state serialization; Added tests for Hadoop Writable type extraction and TypeInformation; Added tests for HybridSource components; Added tests for Predicates utility class; Added tests for Writable serializer and comparator runtime components; Added tests for batch file compaction and partition committing; Added tests for client heartbeat behavior and job initialization error handling; Added tests for stream file compaction operators; Added tests for streaming file writer and partition commit info; Added tests for the new FLIP-27 FileSource implementation; Added tests to verify Scala-free Flink execution; Added unit and integration tests for application mode deployment utilities; Added unit tests for AsyncSinkWriter and supporting components; Added unit tests for CompactorRequestTypeInfo; Added unit tests for ExponentialWaitStrategy; Added unit tests for FileCommitter; Added unit tests for FileSink compactor components; Added unit tests for Flink CLI commands; Added unit tests for FutureCompletingBlockingQueue synchronization; Added unit tests for Hadoop compatibility connectors; Added unit tests for HadoopUtils parameter parsing; Added unit tests for JarHandlerUtils and supporting test programs; Added unit tests for LocalityAwareSplitAssigner; Added unit tests for REST client checkpoint, savepoint, and configuration handling; Added unit tests for SerdeUtils serialization logic; Added unit tests for SplitFetcher and SplitFetcherManager; Added unit tests for file source connector implementation components; Added unit tests for file table connector components; ArchUnit test infrastructure added to flink-formats modules; ArchUnit test rules for Hadoop compatibility test code; ArchUnit test rules for the Files connector updated to support JUnit 5 parallel execution; Centralized architectural test suite for production code; DataGeneratorSource tests moved to dedicated test module; Migrate flink-fs-tests to JUnit 5 and add HDFS integration tests; Migrate flink-tests to JUnit 5; New FileSink end-to-end test program added; New SQL Client end-to-end tests with custom toolbox functions and JUnit 5 migration; New end-to-end tests for SQL table planner features; Updated ArchUnit violation stores for test code rules.

Dependencies

Massive dependency convergence and security updates across 182 manifests

This change updates dependencies across 182 Maven manifests to converge versions, resolve conflicts, and address security vulnerabilities. Key updates include upgrading Netty to 4.2.6.Final (and later 4.2.13.Final), Jackson to 2.15.3, Log4j to 2.25.3, and Commons-Compress to 1.26.0 to address [CVE redacted]. Other notable bumps include Guava to 32.1.3-jre, Mockito to 5.19.0, and the Flink shaded libraries to 21.0. The update also migrates the RPC layer from Akka to Pekko 1.7.0 and updates the Kubernetes client to Fabric8 5.12.4.

(dependencies) · high confidence

Maven wrapper upgraded to version 3.3.4 with Maven 3.9.16 distribution

The Maven wrapper configuration has been updated to use wrapper version 3.3.4 and the Maven 3.9.16 distribution binary. This change includes SHA-256 checksums for both the wrapper JAR and the Maven distribution to ensure integrity during download, replacing previous versions and ensuring consistent build environments across development and CI.

.mvn/wrapper · high confidence

Upgrade Netty to 4.2.6.Final

The Flink runtime has upgraded the Netty networking library to version 4.2.6.Final. This update brings the underlying network stack to a newer version, which may include performance improvements, bug fixes, and compatibility updates for the transport layer used by Flink's internal communication.

flink-runtime · high confidence

Housekeeping

Add NOTICE and license files for the table-planner-loader-bundle; Updated legal notices and license files for bundled dependencies.

Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.

How this codebase got here

Score

  • CAI 49 → 53 (+3.4)
  • Rubric changed (rubric-2026.08.19 → rubric-2026.09.15) — scores are not directly comparable.

Lenses

  • Code Health 85 → 88 (+2.4)
  • Architecture 100 → 94 (-6.1)
  • Maturity 68 → 72 (+3.3)
  • Readiness 36 → 52 (+16.1)
  • Security 72 → 62 (-10.8)
  • Domain Modelling 100 → 92 (-7.6)
  • Accessibility 43 → 41 (-2.0)

Resolved (609)

  • 1.accept (cognitive 46) (flink-yarn-tests/src/test/java/org/apache/flink/yarn/YarnTestBase.java)
  • AnswerFormatter.writeContent (cognitive 16) (flink-end-to-end-tests/flink-tpcds-test/src/main/java/org/apache/flink/table/tpcds/utils/AnswerFormatter.java)
  • AutoRescalingITCase.testCheckpointRescalingPartitionedOperatorState (cognitive 16) (flink-tests/src/test/java/org/apache/flink/test/checkpointing/AutoRescalingITCase.java)
  • Boundary-crossing change coupling: FsStateChangelogWriter.java ↔ ChangelogKeyedStateBackend.java (flink-dstl/flink-dstl-dfs/src/main/java/org/apache/flink/changelog/fs/FsStateChangelogWriter.java)
  • Builder.buildArgumentInfo (cognitive 26) (flink-table/flink-table-test-utils/src/main/java/org/apache/flink/table/runtime/functions/ProcessTableFunctionTestHarness.java)
  • Builder.deriveOutputType (cognitive 32) (flink-table/flink-table-test-utils/src/main/java/org/apache/flink/table/runtime/functions/ProcessTableFunctionTestHarness.java)
  • Builder.validateInitialStateKeys (cognitive 17) (flink-table/flink-table-test-utils/src/main/java/org/apache/flink/table/runtime/functions/ProcessTableFunctionTestHarness.java)
  • Builder.validatePartitionConsistency (cognitive 16) (flink-table/flink-table-test-utils/src/main/java/org/apache/flink/table/runtime/functions/ProcessTableFunctionTestHarness.java)
  • ConnectedComponentsData.getRandomOddEvenEdges (cognitive 16) (flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/testdata/ConnectedComponentsData.java)
  • CountNewsClicksProcessFunction.processRecord (cognitive 16) (flink-examples/flink-examples-streaming/src/main/java/org/apache/flink/streaming/examples/dsv2/eventtime/CountNewsClicks.java)
  • Coverage not included — suite not readable by the collector
  • Dependency hygiene not measured — dependency manifest found but not parsed for hygiene
  • Duplicated block (10 lines × 13) (flink-clients/src/main/java/org/apache/flink/client/deployment/application/PackagedProgramApplicationEntry.java)
  • Duplicated block (10 lines × 2) (flink-clients/src/main/java/org/apache/flink/client/cli/CliFrontendParser.java)
  • Duplicated block (10 lines × 2) (flink-connectors/flink-connector-datagen/src/main/java/org/apache/flink/connector/datagen/functions/FromElementsGeneratorFunction.java)
  • Duplicated block (10 lines × 2) (flink-connectors/flink-connector-files/src/main/java/org/apache/flink/connector/file/table/FileSystemTableSink.java)
  • Duplicated block (10 lines × 2) (flink-core/src/main/java/org/apache/flink/api/common/operators/CollectionExecutor.java)
  • Duplicated block (10 lines × 2) (flink-core/src/main/java/org/apache/flink/api/common/operators/base/CoGroupOperatorBase.java)
  • Duplicated block (10 lines × 2) (flink-core/src/main/java/org/apache/flink/api/common/operators/base/FlatMapOperatorBase.java)
  • Duplicated block (10 lines × 2) (flink-core/src/main/java/org/apache/flink/api/common/typeutils/base/MapSerializer.java)
  • …and 589 more

New (2707)

  • AggCallSelectivityEstimator.estimateComparison (cognitive 22) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/metadata/AggCallSelectivityEstimator.scala)
  • AggCallSelectivityEstimator.estimateComparison (cyclomatic 19) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/metadata/AggCallSelectivityEstimator.scala)
  • AggCallSelectivityEstimator.estimateNumericComparison (cognitive 36) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/metadata/AggCallSelectivityEstimator.scala)
  • AggCallSelectivityEstimator.estimateNumericComparison (cyclomatic 20) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/metadata/AggCallSelectivityEstimator.scala)
  • AggCallSelectivityEstimator.getAggCallInterval (cognitive 22) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/metadata/AggCallSelectivityEstimator.scala)
  • AggFunctionFactory.createAggFunction (cognitive 17) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/AggFunctionFactory.scala)
  • AggFunctionFactory.createAggFunction (cyclomatic 44) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/utils/AggFunctionFactory.scala)
  • AppInterceptor.intercept (cognitive 17) (flink-runtime-web/web-dashboard/src/app/app.interceptor.ts)
  • AppInterceptor.intercept (cyclomatic 21) (flink-runtime-web/web-dashboard/src/app/app.interceptor.ts)
  • BaseAPIStabilityDecorator.call (cognitive 17) (flink-python/pyflink/util/api_stability_decorators.py)
  • BatchPhysicalHashJoinRule.onMatch (cognitive 17) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/physical/batch/BatchPhysicalHashJoinRule.scala)
  • BatchPhysicalJoinBase.satisfyHashDistributionOnNonBroadcastJoin (cognitive 18) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/batch/BatchPhysicalJoinBase.scala)
  • BatchPhysicalJoinRuleBase.checkBroadcast (cyclomatic 18) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/physical/batch/BatchPhysicalJoinRuleBase.scala)
  • BatchPhysicalLegacySinkRule.convert (cognitive 16) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/physical/batch/BatchPhysicalLegacySinkRule.scala)
  • BatchPhysicalOverAggregateBase.satisfyTraits (cognitive 25) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/batch/BatchPhysicalOverAggregateBase.scala)
  • BatchPhysicalOverAggregateRule.onMatch (cognitive 16) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/physical/batch/BatchPhysicalOverAggregateRule.scala)
  • BatchPhysicalPythonGroupAggregate.satisfyTraits (cognitive 16) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/batch/BatchPhysicalPythonGroupAggregate.scala)
  • BatchPhysicalSinkRule.convert (cognitive 25) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/physical/batch/BatchPhysicalSinkRule.scala)
  • BatchPhysicalSortAggregate.satisfyTraits (cognitive 16) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/nodes/physical/batch/BatchPhysicalSortAggregate.scala)
  • BatchPhysicalWindowAggregateRule.transformTimeSlidingWindow (cognitive 23) (flink-table/flink-table-planner/src/main/scala/org/apache/flink/table/planner/plan/rules/physical/batch/BatchPhysicalWindowAggregateRule.scala)
  • …and 2687 more

Changes since last survey

  • 239 commits — 203 feature/other, 36 fixes

By area

  • flink-runtime/src — 37 commits
  • flink-table/flink-table-planner — 37 commits
  • flink-python/pyflink — 24 commits
  • (root) — 13 commits
  • flink-libraries/flink-state-processing-api — 11 commits
  • flink-table/flink-table-common — 11 commits
  • flink-filesystems/flink-s3-fs-native — 10 commits
  • flink-core/src — 8 commits
  • flink-table/flink-table-runtime — 7 commits
  • docs/content — 6 commits
  • flink-python/src — 5 commits
  • flink-runtime-web/web-dashboard — 5 commits
  • flink-tests/src — 5 commits
  • .github/workflows — 4 commits
  • flink-table/flink-table-api-java — 4 commits
  • tools/azure-pipelines — 4 commits
  • flink-end-to-end-tests/test-scripts — 3 commits
  • flink-python/pom.xml — 3 commits
  • flink-runtime-web/src — 3 commits
  • flink-state-backends/flink-statebackend-rocksdb — 3 commits

Notable commits

  • fix: Revert "[FLINK-39977][runtime] Add ITCase for file-merged channel state recovery"
  • fix: Revert "[FLINK-39977][runtime] Recovery of merged channel state handles"
  • fix: [FLINK-23758][python] Fix semantics of Expression.invert (#29148)
  • fix: [FLINK-25802][FLINK-30499][table] Fix TIMESTAMP codegen for RANGE OVER window bounds
  • fix: [FLINK-39014][table] Fix the conversion to relational algebra issue in Batch Mode (#28500)
  • fix: [FLINK-40070][tests] Fix flaky DynamicParameterITCase reading rolled JobManager logs
  • fix: [FLINK-40271][table] Fix SQL serialization of Table API group window properties
  • fix: [FLINK-40300][python] Fix literal expression in Python Table API (#28924)
  • fix: [FLINK-40335][runtime] Fix Hybrid Shuffle index corruption caused by shared ByteBuffer
  • fix: [FLINK-40370][python] Fix bytes-backed Avro decimal framing (#28996)
  • fix: [FLINK-40415][python] Fix Python 3.9 dependency conflict of grpcio-tools (#28992)
  • fix: [FLINK-40462][ci] Fix Build Python Wheels for GHA
  • fix: [FLINK-40475][runtime] Fix watermark loss and stall in StatusWatermarkValve when subpartitions realign after idleness (#29024)
  • fix: [FLINK-40477][table] Fix constraint enforcer for partial deletes (#29067)
  • fix: [FLINK-40481][table] Fix SHOW CREATE FROM_NOW timestamp direction (#29030)
  • fix: [FLINK-40568][tests] Fix flaky ProcessTableFunctionSemanticTests by using materialized data assertion
  • fix: [FLINK-40592][metrics] Fix PushGateway basic authentication without JAXB (#29151)
  • fix: [FLINK-40597][core] Fix input type validation for Types.UUID
  • fix: [FLINK-40606][table-planner] Fix code generation for binary array casts
  • fix: [FLINK-40654][ci] Fix free_disk_space.sh path in build-nightly-dist.yml
  • …and 219 more

Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.

Survey your own repository

apache/flink was measured the same way every project in this corpus was: the same rubric, at a pinned commit, with the result published in full. Point a surveyor at a repository you know and see whether you agree with it.

About this page

  • The score is its most recent published measurement, taken on 24 September 2026 at a pinned commit. It is not a live figure and does not change until the project is measured again.
  • Measured at commit 654f0edc88676df691b8a59f5e44186602b5820f — the exact code this score is about.
  • Scored under rubric-2026.09.15 — the same rubric and the same method as every other entry in this index.
  • Measured by watchdog.canine.dev using codehealth-analyzer preprod-ae95d6cad036.