Skip to content

Commit f6b0452

Browse files
committed
Enhance DataFlowSpec with new features: replace Errors field type with ErrorSinkSpec, add CheckpointReset and StrictIdempotency fields. Update validation and processing logic to accommodate these changes, including new tests for error sink configuration and checkpoint handling.
1 parent da54b8d commit f6b0452

24 files changed

Lines changed: 645 additions & 28 deletions
Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,30 @@
1+
/*
2+
Copyright 2024.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package v1
18+
19+
// AnnotationResetCheckpoint requests a one-shot checkpoint reset on the next processor run.
20+
const AnnotationResetCheckpoint = "dataflow.dataflow.io/reset-checkpoint"
21+
22+
// CheckpointResetRequested reports whether the processor should clear persisted checkpoint on startup.
23+
func CheckpointResetRequested(spec *DataFlowSpec) bool {
24+
return spec != nil && spec.CheckpointReset != nil && *spec.CheckpointReset
25+
}
26+
27+
// StrictIdempotencyEnabled reports whether non-idempotent polling pipelines are rejected at admission.
28+
func StrictIdempotencyEnabled(spec *DataFlowSpec) bool {
29+
return spec != nil && spec.StrictIdempotency != nil && *spec.StrictIdempotency
30+
}

‎api/v1/dataflow_errors.go‎

Lines changed: 63 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,63 @@
1+
/*
2+
Copyright 2024.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package v1
18+
19+
const (
20+
// ErrorAckPolicyAfterWrite commits source offsets after the error sink successfully writes (default).
21+
ErrorAckPolicyAfterWrite = "afterWrite"
22+
// ErrorAckPolicyNever does not commit source offsets for messages sent to the error sink.
23+
ErrorAckPolicyNever = "never"
24+
// ErrorAckPolicyAfterMainSinkSuccess commits offsets only after main sink success (same as never on error path).
25+
ErrorAckPolicyAfterMainSinkSuccess = "afterMainSinkSuccess"
26+
)
27+
28+
// ErrorSinkSpec defines the error sink and how source ack behaves for failed messages.
29+
type ErrorSinkSpec struct {
30+
SinkSpec `json:",inline"`
31+
32+
// AckPolicy controls when source offsets/checkpoints are committed for messages routed to the error sink.
33+
// afterWrite (default): ack after error sink write.
34+
// never: do not ack failed messages (they may be re-read on restart).
35+
// afterMainSinkSuccess: ack only after main sink success (failed messages are not acked).
36+
// +kubebuilder:validation:Enum=afterWrite;never;afterMainSinkSuccess
37+
// +kubebuilder:default:=afterWrite
38+
// +optional
39+
AckPolicy string `json:"ackPolicy,omitempty"`
40+
}
41+
42+
// ErrorAckPolicyOrDefault returns the configured error sink ack policy.
43+
func ErrorAckPolicyOrDefault(errors *ErrorSinkSpec) string {
44+
if errors == nil || errors.AckPolicy == "" {
45+
return ErrorAckPolicyAfterWrite
46+
}
47+
switch errors.AckPolicy {
48+
case ErrorAckPolicyNever, ErrorAckPolicyAfterMainSinkSuccess:
49+
return errors.AckPolicy
50+
default:
51+
return ErrorAckPolicyAfterWrite
52+
}
53+
}
54+
55+
// ShouldAckOnErrorSink reports whether source ack should propagate when writing to the error sink.
56+
func ShouldAckOnErrorSink(errors *ErrorSinkSpec) bool {
57+
switch ErrorAckPolicyOrDefault(errors) {
58+
case ErrorAckPolicyNever, ErrorAckPolicyAfterMainSinkSuccess:
59+
return false
60+
default:
61+
return true
62+
}
63+
}

‎api/v1/dataflow_errors_test.go‎

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
/*
2+
Copyright 2024.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package v1
18+
19+
import (
20+
"testing"
21+
22+
"github.com/stretchr/testify/assert"
23+
)
24+
25+
func TestShouldAckOnErrorSink(t *testing.T) {
26+
t.Parallel()
27+
28+
assert.True(t, ShouldAckOnErrorSink(nil))
29+
assert.True(t, ShouldAckOnErrorSink(&ErrorSinkSpec{}))
30+
assert.True(t, ShouldAckOnErrorSink(&ErrorSinkSpec{AckPolicy: ErrorAckPolicyAfterWrite}))
31+
assert.False(t, ShouldAckOnErrorSink(&ErrorSinkSpec{AckPolicy: ErrorAckPolicyNever}))
32+
assert.False(t, ShouldAckOnErrorSink(&ErrorSinkSpec{AckPolicy: ErrorAckPolicyAfterMainSinkSuccess}))
33+
}
34+
35+
func TestCheckpointResetRequested(t *testing.T) {
36+
t.Parallel()
37+
38+
assert.False(t, CheckpointResetRequested(nil))
39+
assert.False(t, CheckpointResetRequested(&DataFlowSpec{}))
40+
trueVal := true
41+
assert.True(t, CheckpointResetRequested(&DataFlowSpec{CheckpointReset: &trueVal}))
42+
}

‎api/v1/dataflow_idempotency.go‎

Lines changed: 105 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,105 @@
1+
/*
2+
Copyright 2024.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package v1
18+
19+
import (
20+
"github.com/dataflow-operator/dataflow/pkg/providers"
21+
"k8s.io/apimachinery/pkg/util/validation/field"
22+
"sigs.k8s.io/controller-runtime/pkg/webhook/admission"
23+
)
24+
25+
func validateErrors(errors *ErrorSinkSpec, f *field.Path) field.ErrorList {
26+
var all field.ErrorList
27+
if errors == nil {
28+
return all
29+
}
30+
all = append(all, validateSink(&errors.SinkSpec, f)...)
31+
if errors.AckPolicy != "" &&
32+
errors.AckPolicy != ErrorAckPolicyAfterWrite &&
33+
errors.AckPolicy != ErrorAckPolicyNever &&
34+
errors.AckPolicy != ErrorAckPolicyAfterMainSinkSuccess {
35+
all = append(all, field.NotSupported(f.Child("ackPolicy"), errors.AckPolicy, []string{
36+
ErrorAckPolicyAfterWrite,
37+
ErrorAckPolicyNever,
38+
ErrorAckPolicyAfterMainSinkSuccess,
39+
}))
40+
}
41+
return all
42+
}
43+
44+
func validateIdempotency(spec *DataFlowSpec, f *field.Path) field.ErrorList {
45+
var all field.ErrorList
46+
if spec == nil || !isPollingSourceType(spec.Source.Type) {
47+
return all
48+
}
49+
if sinkIsIdempotent(&spec.Sink) {
50+
return all
51+
}
52+
msg := "polling source with non-idempotent main sink may produce duplicates on restart; enable upsertMode on the sink or set strictIdempotency: false to accept this warning"
53+
if StrictIdempotencyEnabled(spec) {
54+
all = append(all, field.Invalid(f.Child("sink"), spec.Sink.Type,
55+
"strictIdempotency is enabled but main sink is not idempotent (enable upsertMode for postgresql/trino/clickhouse sinks)"))
56+
} else {
57+
_ = msg // warnings emitted via WarnDataFlowSpec
58+
}
59+
return all
60+
}
61+
62+
// WarnDataFlowSpec returns admission warnings for non-fatal spec issues.
63+
func WarnDataFlowSpec(spec *DataFlowSpec) admission.Warnings {
64+
var warnings admission.Warnings
65+
if spec == nil {
66+
return warnings
67+
}
68+
if isPollingSourceType(spec.Source.Type) && !sinkIsIdempotent(&spec.Sink) && !StrictIdempotencyEnabled(spec) {
69+
warnings = append(warnings,
70+
"polling source with non-idempotent main sink may produce duplicates on restart; enable sink.config.upsertMode or set strictIdempotency: true to reject at admission")
71+
}
72+
return warnings
73+
}
74+
75+
func isPollingSourceType(sourceType string) bool {
76+
return providers.SourceSupportsCheckpoint(sourceType)
77+
}
78+
79+
func sinkIsIdempotent(sink *SinkSpec) bool {
80+
if sink == nil || sink.Config == nil || len(sink.Config.Raw) == 0 {
81+
return false
82+
}
83+
switch sink.Type {
84+
case "postgresql":
85+
cfg, err := sink.GetPostgreSQLConfig()
86+
if err != nil || cfg == nil {
87+
return false
88+
}
89+
return cfg.UpsertMode != nil && *cfg.UpsertMode
90+
case "clickhouse":
91+
cfg, err := sink.GetClickHouseConfig()
92+
if err != nil || cfg == nil {
93+
return false
94+
}
95+
return cfg.UpsertMode != nil && *cfg.UpsertMode
96+
case "trino":
97+
cfg, err := sink.GetTrinoConfig()
98+
if err != nil || cfg == nil {
99+
return false
100+
}
101+
return cfg.UpsertMode != nil && *cfg.UpsertMode
102+
default:
103+
return false
104+
}
105+
}
Lines changed: 69 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,69 @@
1+
/*
2+
Copyright 2024.
3+
4+
Licensed under the Apache License, Version 2.0 (the "License");
5+
you may not use this file except in compliance with the License.
6+
You may obtain a copy of the License at
7+
8+
http://www.apache.org/licenses/LICENSE-2.0
9+
10+
Unless required by applicable law or agreed to in writing, software
11+
distributed under the License is distributed on an "AS IS" BASIS,
12+
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13+
See the License for the specific language governing permissions and
14+
limitations under the License.
15+
*/
16+
17+
package v1
18+
19+
import (
20+
"testing"
21+
22+
"github.com/stretchr/testify/assert"
23+
)
24+
25+
func TestWarnDataFlowSpec_pollingNonIdempotent(t *testing.T) {
26+
t.Parallel()
27+
28+
upsertOff := false
29+
spec := DataFlowSpec{
30+
Source: SourceSpec{Type: "postgresql", Config: mustConfig(PostgreSQLSourceSpec{ConnectionString: "postgres://u:p@localhost/db", Table: "t"})},
31+
Sink: SinkSpec{
32+
Type: "postgresql",
33+
Config: mustConfig(PostgreSQLSinkSpec{ConnectionString: "postgres://u:p@localhost/db", Table: "out", UpsertMode: &upsertOff}),
34+
},
35+
}
36+
warnings := WarnDataFlowSpec(&spec)
37+
assert.Len(t, warnings, 1)
38+
}
39+
40+
func TestValidateDataFlowSpec_strictIdempotency(t *testing.T) {
41+
t.Parallel()
42+
43+
strict := true
44+
spec := DataFlowSpec{
45+
StrictIdempotency: &strict,
46+
Source: SourceSpec{Type: "postgresql", Config: mustConfig(PostgreSQLSourceSpec{ConnectionString: "postgres://u:p@localhost/db", Table: "t"})},
47+
Sink: SinkSpec{
48+
Type: "postgresql",
49+
Config: mustConfig(PostgreSQLSinkSpec{ConnectionString: "postgres://u:p@localhost/db", Table: "out"}),
50+
},
51+
}
52+
errs := ValidateDataFlowSpec(&spec)
53+
assert.NotEmpty(t, errs)
54+
}
55+
56+
func TestValidateErrors_ackPolicy(t *testing.T) {
57+
t.Parallel()
58+
59+
spec := DataFlowSpec{
60+
Source: SourceSpec{Type: "kafka", Config: mustConfig(KafkaSourceSpec{Brokers: []string{"b"}, Topic: "t"})},
61+
Sink: SinkSpec{Type: "kafka", Config: mustConfig(KafkaSinkSpec{Brokers: []string{"b"}, Topic: "out"})},
62+
Errors: &ErrorSinkSpec{
63+
SinkSpec: SinkSpec{Type: "kafka", Config: mustConfig(KafkaSinkSpec{Brokers: []string{"b"}, Topic: "err"})},
64+
AckPolicy: "invalid",
65+
},
66+
}
67+
errs := ValidateDataFlowSpec(&spec)
68+
assert.NotEmpty(t, errs)
69+
}

‎api/v1/dataflow_types.go‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,7 @@ type DataFlowSpec struct {
3838

3939
// Errors defines the error sink for messages that failed to be written to the main sink
4040
// +optional
41-
Errors *SinkSpec `json:"errors,omitempty"`
41+
Errors *ErrorSinkSpec `json:"errors,omitempty"`
4242

4343
// Resources defines the resource requirements for the processor pod
4444
// +optional
@@ -96,6 +96,16 @@ type DataFlowSpec struct {
9696
// +optional
9797
AckGranularity string `json:"ackGranularity,omitempty"`
9898

99+
// CheckpointReset clears persisted source checkpoint on the next processor start (one-shot).
100+
// Alternatively set annotation dataflow.dataflow.io/reset-checkpoint: "true" on the DataFlow.
101+
// +optional
102+
CheckpointReset *bool `json:"checkpointReset,omitempty"`
103+
104+
// StrictIdempotency rejects polling sources paired with non-idempotent main sinks at admission.
105+
// When false (default), a warning is emitted instead.
106+
// +optional
107+
StrictIdempotency *bool `json:"strictIdempotency,omitempty"`
108+
99109
// ChannelBufferSize is the buffer size for message channels between source, processor, and sink.
100110
// Larger values reduce blocking when sink is slower than source (e.g. high Kafka throughput).
101111
// Default: 100. Recommended for high throughput: 500–1000.

‎api/v1/dataflow_validation.go‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,12 +40,13 @@ func ValidateDataFlowSpec(spec *DataFlowSpec) field.ErrorList {
4040
all = append(all, validateSource(&spec.Source, f.Child("source"))...)
4141
all = append(all, validateSink(&spec.Sink, f.Child("sink"))...)
4242
if spec.Errors != nil {
43-
all = append(all, validateSink(spec.Errors, f.Child("errors"))...)
43+
all = append(all, validateErrors(spec.Errors, f.Child("errors"))...)
4444
}
4545
all = append(all, validateTransformations(spec.Transformations, f.Child("transformations"))...)
4646
all = append(all, validateResources(spec.Resources, f.Child("resources"))...)
4747
all = append(all, validateReplicas(spec, f)...)
4848
all = append(all, validateAckGranularity(spec, f)...)
49+
all = append(all, validateIdempotency(spec, f)...)
4950
return all
5051
}
5152

‎api/v1/dataflow_webhook.go‎

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -30,10 +30,10 @@ func (r *DataFlow) ValidateCreate(ctx context.Context, obj runtime.Object) (admi
3030
return nil, nil
3131
}
3232
errs := ValidateDataFlowSpec(&df.Spec)
33-
if len(errs) == 0 {
34-
return nil, nil
33+
if len(errs) > 0 {
34+
return nil, errs.ToAggregate()
3535
}
36-
return nil, errs.ToAggregate()
36+
return WarnDataFlowSpec(&df.Spec), nil
3737
}
3838

3939
// ValidateUpdate implements admission.CustomValidator.
@@ -43,10 +43,10 @@ func (r *DataFlow) ValidateUpdate(ctx context.Context, oldObj, newObj runtime.Ob
4343
return nil, nil
4444
}
4545
errs := ValidateDataFlowSpec(&df.Spec)
46-
if len(errs) == 0 {
47-
return nil, nil
46+
if len(errs) > 0 {
47+
return nil, errs.ToAggregate()
4848
}
49-
return nil, errs.ToAggregate()
49+
return WarnDataFlowSpec(&df.Spec), nil
5050
}
5151

5252
// ValidateDelete implements admission.CustomValidator.

0 commit comments

Comments
 (0)