Skip to content

Commit b19ed21

Browse files
committed
refatoring
1 parent c7effe4 commit b19ed21

31 files changed

Lines changed: 1171 additions & 119 deletions

‎api/v1/Untitled‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
protocol

‎api/v1/dataflow_checkpoint.go‎

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,48 @@
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+
"time"
21+
22+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
23+
)
24+
25+
const defaultCheckpointSaveInterval = 30 * time.Second
26+
27+
// CheckpointPersistenceEnabled reports whether checkpoint persistence is enabled (default true).
28+
func CheckpointPersistenceEnabled(spec *DataFlowSpec) bool {
29+
return spec == nil || spec.CheckpointPersistence == nil || *spec.CheckpointPersistence
30+
}
31+
32+
// CheckpointSyncOnAckEnabled reports whether checkpoint should flush after each sink batch ack.
33+
func CheckpointSyncOnAckEnabled(spec *DataFlowSpec) bool {
34+
return spec != nil && spec.CheckpointSyncOnAck != nil && *spec.CheckpointSyncOnAck
35+
}
36+
37+
// CheckpointSaveIntervalDuration returns the debounce/sync coalesce interval (default 30s).
38+
func CheckpointSaveIntervalDuration(spec *DataFlowSpec) time.Duration {
39+
if spec != nil && spec.CheckpointSaveInterval != nil && spec.CheckpointSaveInterval.Duration > 0 {
40+
return spec.CheckpointSaveInterval.Duration
41+
}
42+
return defaultCheckpointSaveInterval
43+
}
44+
45+
// CheckpointSaveIntervalOrDefault returns the spec interval pointer or a default metav1.Duration.
46+
func CheckpointSaveIntervalOrDefault(spec *DataFlowSpec) metav1.Duration {
47+
return metav1.Duration{Duration: CheckpointSaveIntervalDuration(spec)}
48+
}

‎api/v1/dataflow_checkpoint_test.go‎

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
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+
"time"
22+
23+
"github.com/stretchr/testify/assert"
24+
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
25+
)
26+
27+
func TestCheckpointPersistenceEnabled(t *testing.T) {
28+
t.Parallel()
29+
30+
trueVal := true
31+
falseVal := false
32+
assert.True(t, CheckpointPersistenceEnabled(nil))
33+
assert.True(t, CheckpointPersistenceEnabled(&DataFlowSpec{}))
34+
assert.True(t, CheckpointPersistenceEnabled(&DataFlowSpec{CheckpointPersistence: &trueVal}))
35+
assert.False(t, CheckpointPersistenceEnabled(&DataFlowSpec{CheckpointPersistence: &falseVal}))
36+
}
37+
38+
func TestCheckpointSyncOnAckEnabled(t *testing.T) {
39+
t.Parallel()
40+
41+
trueVal := true
42+
assert.False(t, CheckpointSyncOnAckEnabled(nil))
43+
assert.False(t, CheckpointSyncOnAckEnabled(&DataFlowSpec{}))
44+
assert.True(t, CheckpointSyncOnAckEnabled(&DataFlowSpec{CheckpointSyncOnAck: &trueVal}))
45+
}
46+
47+
func TestCheckpointSaveIntervalDuration(t *testing.T) {
48+
t.Parallel()
49+
50+
assert.Equal(t, 30*time.Second, CheckpointSaveIntervalDuration(nil))
51+
assert.Equal(t, 5*time.Second, CheckpointSaveIntervalDuration(&DataFlowSpec{
52+
CheckpointSaveInterval: &metav1.Duration{Duration: 5 * time.Second},
53+
}))
54+
}

‎api/v1/dataflow_types.go‎

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,18 @@ type DataFlowSpec struct {
7676
// +optional
7777
CheckpointPersistence *bool `json:"checkpointPersistence,omitempty"`
7878

79+
// CheckpointSyncOnAck persists checkpoint to the ConfigMap immediately after each sink batch ack.
80+
// Shrinks the duplicate window on pod crash from the debounce interval to roughly one batch.
81+
// Default: false (debounced save only). Recommended for migration and cron workloads.
82+
// +optional
83+
CheckpointSyncOnAck *bool `json:"checkpointSyncOnAck,omitempty"`
84+
85+
// CheckpointSaveInterval controls how often pending checkpoints are flushed to the ConfigMap.
86+
// Also used as the minimum coalesce interval when checkpointSyncOnAck is true.
87+
// Default: 30s.
88+
// +optional
89+
CheckpointSaveInterval *metav1.Duration `json:"checkpointSaveInterval,omitempty"`
90+
7991
// ChannelBufferSize is the buffer size for message channels between source, processor, and sink.
8092
// Larger values reduce blocking when sink is slower than source (e.g. high Kafka throughput).
8193
// Default: 100. Recommended for high throughput: 500–1000.
@@ -733,6 +745,17 @@ type PostgreSQLSinkSpec struct {
733745
// +optional
734746
ConflictKey *string `json:"conflictKey,omitempty"`
735747

748+
// UpsertVersionColumn is the column used to compare row versions when UpsertStrategy is ifNewer.
749+
// On conflict, the existing row is updated only when EXCLUDED.<column> > target.<column>.
750+
// +optional
751+
UpsertVersionColumn *string `json:"upsertVersionColumn,omitempty"`
752+
753+
// UpsertStrategy controls conflict resolution: always (default) updates on every conflict;
754+
// ifNewer updates only when the incoming version is newer (requires upsertVersionColumn).
755+
// +optional
756+
// +kubebuilder:validation:Enum=always;ifNewer
757+
UpsertStrategy *string `json:"upsertStrategy,omitempty"`
758+
736759
// SoftDeleteColumn specifies column for soft delete (e.g. "deleted_at"). If set, DELETE operations will UPDATE this column instead of physical delete.
737760
// +optional
738761
SoftDeleteColumn *string `json:"softDeleteColumn,omitempty"`
@@ -785,6 +808,16 @@ type TrinoSinkSpec struct {
785808
// +optional
786809
AutoCreateTable *bool `json:"autoCreateTable,omitempty"`
787810

811+
// UpsertMode enables idempotent writes via MERGE (Iceberg catalog).
812+
// Requires conflictKey. Retries and replays update existing rows instead of creating duplicates.
813+
// +optional
814+
UpsertMode *bool `json:"upsertMode,omitempty"`
815+
816+
// ConflictKey specifies the column used to match rows in MERGE ON clause.
817+
// Required when upsertMode is true.
818+
// +optional
819+
ConflictKey *string `json:"conflictKey,omitempty"`
820+
788821
// RawMode when true, creates table with data VARCHAR column (JSON storage).
789822
// When false (default), uses columnar format matching message keys to table columns.
790823
// +optional
@@ -833,6 +866,27 @@ type ClickHouseSinkSpec struct {
833866
// +optional
834867
AutoCreateTable *bool `json:"autoCreateTable,omitempty"`
835868

869+
// UpsertMode enables idempotent writes via ReplacingMergeTree deduplication.
870+
// Inserts remain INSERT; duplicates are resolved on merge using conflictKey ORDER BY.
871+
// +optional
872+
UpsertMode *bool `json:"upsertMode,omitempty"`
873+
874+
// ConflictKey specifies the deduplication key column for ORDER BY when upsertMode is true.
875+
// If not set, the first message column (or created_at in raw mode) is used.
876+
// +optional
877+
ConflictKey *string `json:"conflictKey,omitempty"`
878+
879+
// TableEngine selects the MergeTree engine variant (default: MergeTree).
880+
// When upsertMode is true and tableEngine is unset, ReplacingMergeTree is used for auto-created tables.
881+
// +optional
882+
// +kubebuilder:validation:Enum=MergeTree;ReplacingMergeTree
883+
TableEngine *string `json:"tableEngine,omitempty"`
884+
885+
// ReplacingVersionColumn is the version column for ReplacingMergeTree(engine).
886+
// Rows with the highest version value are kept during background merges.
887+
// +optional
888+
ReplacingVersionColumn *string `json:"replacingVersionColumn,omitempty"`
889+
836890
// RawMode when true, creates table with data String and created_at columns (JSON storage).
837891
// When false (default), creates table from message structure (columnar, replicates source schema).
838892
// +optional

‎api/v1/dataflow_validation.go‎

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -479,6 +479,26 @@ func validatePostgreSQLSink(p *PostgreSQLSinkSpec, f *field.Path) field.ErrorLis
479479
if p.TableSecretRef != nil {
480480
all = append(all, validateSecretRef(p.TableSecretRef, f.Child("tableSecretRef"))...)
481481
}
482+
if p.UpsertMode != nil && *p.UpsertMode {
483+
if err := validateSQLIdentifier(resolveConflictKeyForValidation(p.ConflictKey), f.Child("conflictKey")); err != nil {
484+
all = append(all, err)
485+
}
486+
}
487+
if p.UpsertStrategy != nil && *p.UpsertStrategy == "ifNewer" {
488+
if p.UpsertVersionColumn == nil || strings.TrimSpace(*p.UpsertVersionColumn) == "" {
489+
all = append(all, field.Required(f.Child("upsertVersionColumn"), "upsertVersionColumn is required when upsertStrategy is ifNewer"))
490+
} else if err := validateSQLIdentifier(*p.UpsertVersionColumn, f.Child("upsertVersionColumn")); err != nil {
491+
all = append(all, err)
492+
}
493+
}
494+
if p.UpsertVersionColumn != nil && strings.TrimSpace(*p.UpsertVersionColumn) != "" {
495+
if err := validateSQLIdentifier(*p.UpsertVersionColumn, f.Child("upsertVersionColumn")); err != nil {
496+
all = append(all, err)
497+
}
498+
}
499+
if p.UpsertStrategy != nil && *p.UpsertStrategy != "" && *p.UpsertStrategy != "always" && *p.UpsertStrategy != "ifNewer" {
500+
all = append(all, field.Invalid(f.Child("upsertStrategy"), *p.UpsertStrategy, "must be always or ifNewer"))
501+
}
482502
all = append(all, validateFlattenMetadataSpec(p.RawMode, p.FlattenMetadataColumns, p.FlattenMetadataColumnsPrefix, f)...)
483503
return all
484504
}
@@ -513,6 +533,19 @@ func validateTrinoSink(t *TrinoSinkSpec, f *field.Path) field.ErrorList {
513533
if t.TableSecretRef != nil {
514534
all = append(all, validateSecretRef(t.TableSecretRef, f.Child("tableSecretRef"))...)
515535
}
536+
if t.UpsertMode != nil && *t.UpsertMode {
537+
conflictKey := resolveConflictKeyForValidation(t.ConflictKey)
538+
if conflictKey == "" {
539+
all = append(all, field.Required(f.Child("conflictKey"), "conflictKey is required when upsertMode is enabled"))
540+
} else if err := validateSQLIdentifier(conflictKey, f.Child("conflictKey")); err != nil {
541+
all = append(all, err)
542+
}
543+
catalog := strings.ToLower(t.Catalog)
544+
if catalog != "" && !strings.Contains(catalog, "iceberg") {
545+
all = append(all, field.Invalid(f.Child("catalog"), t.Catalog,
546+
"upsertMode requires an Iceberg catalog (catalog name should contain \"iceberg\")"))
547+
}
548+
}
516549
all = append(all, validateFlattenMetadataSpec(t.RawMode, t.FlattenMetadataColumns, t.FlattenMetadataColumnsPrefix, f)...)
517550
return all
518551
}
@@ -558,10 +591,32 @@ func validateClickHouseSink(c *ClickHouseSinkSpec, f *field.Path) field.ErrorLis
558591
if c.TableSecretRef != nil {
559592
all = append(all, validateSecretRef(c.TableSecretRef, f.Child("tableSecretRef"))...)
560593
}
594+
if c.UpsertMode != nil && *c.UpsertMode {
595+
if c.ConflictKey != nil && strings.TrimSpace(*c.ConflictKey) != "" {
596+
if err := validateSQLIdentifier(*c.ConflictKey, f.Child("conflictKey")); err != nil {
597+
all = append(all, err)
598+
}
599+
}
600+
}
601+
if c.TableEngine != nil && *c.TableEngine != "" && *c.TableEngine != "MergeTree" && *c.TableEngine != "ReplacingMergeTree" {
602+
all = append(all, field.Invalid(f.Child("tableEngine"), *c.TableEngine, "must be MergeTree or ReplacingMergeTree"))
603+
}
604+
if c.ReplacingVersionColumn != nil && strings.TrimSpace(*c.ReplacingVersionColumn) != "" {
605+
if err := validateSQLIdentifier(*c.ReplacingVersionColumn, f.Child("replacingVersionColumn")); err != nil {
606+
all = append(all, err)
607+
}
608+
}
561609
all = append(all, validateFlattenMetadataSpec(c.RawMode, c.FlattenMetadataColumns, c.FlattenMetadataColumnsPrefix, f)...)
562610
return all
563611
}
564612

613+
func resolveConflictKeyForValidation(conflictKey *string) string {
614+
if conflictKey == nil {
615+
return ""
616+
}
617+
return strings.TrimSpace(*conflictKey)
618+
}
619+
565620
func validateSecretRef(r *SecretRef, f *field.Path) field.ErrorList {
566621
var all field.ErrorList
567622
if r == nil {
Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,76 @@
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+
"k8s.io/apimachinery/pkg/util/validation/field"
24+
)
25+
26+
func TestValidatePostgreSQLSink_ifNewerRequiresVersionColumn(t *testing.T) {
27+
strategy := "ifNewer"
28+
spec := &PostgreSQLSinkSpec{
29+
ConnectionString: "postgres://localhost/db",
30+
Table: "t",
31+
UpsertStrategy: &strategy,
32+
}
33+
errs := validatePostgreSQLSink(spec, field.NewPath("sink"))
34+
assert.NotEmpty(t, errs)
35+
}
36+
37+
func TestValidateTrinoSink_upsertRequiresIcebergCatalog(t *testing.T) {
38+
upsertMode := true
39+
conflictKey := "id"
40+
spec := &TrinoSinkSpec{
41+
ServerURL: "http://trino:8080",
42+
Catalog: "hive",
43+
Schema: "default",
44+
Table: "t",
45+
UpsertMode: &upsertMode,
46+
ConflictKey: &conflictKey,
47+
}
48+
errs := validateTrinoSink(spec, field.NewPath("sink"))
49+
assert.NotEmpty(t, errs)
50+
}
51+
52+
func TestValidateTrinoSink_upsertValidIceberg(t *testing.T) {
53+
upsertMode := true
54+
conflictKey := "id"
55+
spec := &TrinoSinkSpec{
56+
ServerURL: "http://trino:8080",
57+
Catalog: "nessie_iceberg",
58+
Schema: "default",
59+
Table: "t",
60+
UpsertMode: &upsertMode,
61+
ConflictKey: &conflictKey,
62+
}
63+
errs := validateTrinoSink(spec, field.NewPath("sink"))
64+
assert.Empty(t, errs)
65+
}
66+
67+
func TestValidateClickHouseSink_upsertEngineFields(t *testing.T) {
68+
badEngine := "SummingMergeTree"
69+
spec := &ClickHouseSinkSpec{
70+
ConnectionString: "clickhouse://localhost:9000",
71+
Table: "t",
72+
TableEngine: &badEngine,
73+
}
74+
errs := validateClickHouseSink(spec, field.NewPath("sink"))
75+
assert.NotEmpty(t, errs)
76+
}

0 commit comments

Comments
 (0)