[STIP-15] [Feature][Core] Design of Dirty Data Collection
- Dominant language
- Java
- Stars
- 9.7k
- Forks
- 2.4k
- Avg merge
- 3d 17h
- Merged PRs (30d)
- 210
Description
### Search before asking
- [X] I had searched in the [feature](https://github.com/apache/incubator-seatunnel/issues?q=is%3Aissue+label%3A%22Feature%22) and found no similar feature requirement.
### Description
# Design of Dirty Data Collection Functionality for Apache SeaTunnel
## Introduction
The dirty data collection function can effectively guide users to discover data quality problems in a timely manner in the data integration framework, ensuring the accuracy and stability of data synchronization.
## Functional Requirements
The dirty data collection function needs to support the following features:
1. Collect and record all dirty data content and exception information.
2. Support plug-in extension.
## Technical Solution
### Interface Design
Add interface `DirtyRecordCollector`:
```java
/** Base interface for dirty records collector */
public interface DirtyRecordCollector extends Serializable {
/** Collect dirty record with exception and error message */
void collect(
final int subTaskIndex,
final SeaTunnelRow dirtyRecord,
final Throwable exception,
final String errorMessage);
/** Collect dirty record with exception */
default void collect(
final int subTaskIndex, final SeaTunnelRow dirtyRecord, final Throwable exception) {
collect(subTaskIndex, dirtyRecord, exception, "");
}
}
```
Add method `getDirtyRecordCollector` to `SinkWriter.Context`:
```java
interface Context extends Serializable {
/** @return The index of this subtask. */
int getIndexOfSubtask();
/** @return metricsContext of this reader. */
MetricsContext getMetricsContext();
/**
* Get dirty record collector
*/
DirtyRecordCollector getDirtyRecordCollector();
}
```
Add method `getPluginConfig` to `SeaTunnelPluginLifeCycle`:
```java
/**
* Get the config of plugin
*/
Config getPluginConfig();
```
### Configuration File Example
```hocon
env {
job.mode = BATCH
dirty.collector = {
type = log
}
}
source {
FakeSource {
row.num = 100
schema {
fields {
name = string
age = int
}
}
}
}
sink {
Console {}
}
```
### Technical Implementation
1. Modify `seatunnel-config-shade`. During the transmission process, `SeaTunnelSink` will undergo serialization operations, but the underlying interface of `typesafe-config` does not support serialization, so it is necessary to make `Config` support serialization to ensure that it is not lost during transmission. See #4586 for specific implementation details.
2. Modify all `SeaTunnelSink` to implement the `getPluginConfig` interface and inject the pluginConfig during the prepare phase.
Before instantiating the Sink, merge the plugin configuration for all Sinks, and merge the dirty data collector configuration information from the env section into the plugin configuration.
3. Modify `DefaultSinkWriterContext` and `SinkWriterContext` to instantiate the `DirtyRecordCollector` using the configuration and implement the `getDirtyRecordCollector` method.
4. Modify all `SeaTunnelSink` to report and update metrics using the dirty data collector when data writing or conversion fails.
### Process Design

### Usage Scenario
_No response_
### Related issues
_No response_
### Are you willing to submit a PR?
- [X] Yes I am willing to submit a PR!
### Code of Conduct
- [X] I agree to follow this project's [Code of Conduct](https://www.apache.org/foundation/policies/conduct)
Contributor guide
No contributing guide indexed for this repository
Research direction
Start by reading the existing SinkWriter.Context, SeaTunnelPluginLifeCycle, DefaultSinkWriterContext, and SinkWriterContext APIs, then review #4586 for the seatunnel-config-shade dependency. Trace how SeaTunnelSink instances receive configuration and handle write or conversion failures. Done means the collector is configurable and extensible, exposed through sink context, and used for dirty records and related metrics across sinks.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java
- Domain
- backend, data-engineering
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Quiet
- Clarity
- Mostly clear
- Newbie friendliness
- 38/100