Skip to content

Commit ea293ad

Browse files
committed
Improve pg connector
1 parent a255335 commit ea293ad

15 files changed

Lines changed: 1557 additions & 338 deletions

‎api/v1/dataflow_types.go‎

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -215,6 +215,20 @@ type PostgreSQLSourceSpec struct {
215215
// RawMode when true, wraps each row as JSON: {"value": <row data>, "_metadata": {table, id}}
216216
// +optional
217217
RawMode *bool `json:"rawMode,omitempty"`
218+
219+
// ReadBatchSize limits rows per poll to reduce DB load (0 = no limit)
220+
// +optional
221+
ReadBatchSize *int32 `json:"readBatchSize,omitempty"`
222+
223+
// ChangeTrackingColumn is the column used to track changes (default: updated_at).
224+
// Not used when Query is specified.
225+
// +optional
226+
// +kubebuilder:default="updated_at"
227+
ChangeTrackingColumn string `json:"changeTrackingColumn,omitempty"`
228+
229+
// AutoCreateTable creates the table if it doesn't exist before reading
230+
// +optional
231+
AutoCreateTable *bool `json:"autoCreateTable,omitempty"`
218232
}
219233

220234
// TrinoSourceSpec defines Trino source configuration
@@ -525,8 +539,9 @@ type PostgreSQLSinkSpec struct {
525539
// Table to write to
526540
Table string `json:"table"`
527541

528-
// BatchSize for batch inserts
542+
// BatchSize for batch inserts (default: 1)
529543
// +optional
544+
// +kubebuilder:default=1
530545
BatchSize *int32 `json:"batchSize,omitempty"`
531546

532547
// BatchFlushIntervalSeconds flushes the batch after this many seconds even if BatchSize is not reached (default: 10; 0 disables timer).
@@ -547,6 +562,15 @@ type PostgreSQLSinkSpec struct {
547562
// +optional
548563
ConflictKey *string `json:"conflictKey,omitempty"`
549564

565+
// SoftDeleteColumn specifies column for soft delete (e.g. "deleted_at"). If set, DELETE operations will UPDATE this column instead of physical delete.
566+
// +optional
567+
SoftDeleteColumn *string `json:"softDeleteColumn,omitempty"`
568+
569+
// RawMode when true, expects messages in format {"value": <data>, "_metadata": {...}}. Table is created with value JSONB and _metadata JSONB columns.
570+
// When false, table structure is inferred from the first message (replicates source structure).
571+
// +optional
572+
RawMode *bool `json:"rawMode,omitempty"`
573+
550574
// ConnectionStringSecretRef references a Kubernetes secret for connection string
551575
// +optional
552576
ConnectionStringSecretRef *SecretRef `json:"connectionStringSecretRef,omitempty"`

‎api/v1/zz_generated.deepcopy.go‎

Lines changed: 20 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎config/crd/bases/dataflow.dataflow.io_dataflows.yaml‎

Lines changed: 84 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -1371,7 +1371,8 @@ spec:
13711371
format: int32
13721372
type: integer
13731373
batchSize:
1374-
description: BatchSize for batch inserts
1374+
default: 1
1375+
description: 'BatchSize for batch inserts (default: 1)'
13751376
format: int32
13761377
type: integer
13771378
conflictKey:
@@ -1400,6 +1401,16 @@ spec:
14001401
- key
14011402
- name
14021403
type: object
1404+
rawMode:
1405+
description: |-
1406+
RawMode when true, expects messages in format {"value": <data>, "_metadata": {...}}. Table is created with value JSONB and _metadata JSONB columns.
1407+
When false, table structure is inferred from the first message (replicates source structure).
1408+
type: boolean
1409+
softDeleteColumn:
1410+
description: SoftDeleteColumn specifies column for soft delete
1411+
(e.g. "deleted_at"). If set, DELETE operations will UPDATE
1412+
this column instead of physical delete.
1413+
type: string
14031414
table:
14041415
description: Table to write to
14051416
type: string
@@ -1705,12 +1716,43 @@ spec:
17051716
required:
17061717
- type
17071718
type: object
1719+
imagePullSecrets:
1720+
description: ImagePullSecrets is a list of references to secrets in
1721+
the same namespace to use for pulling the processor image from a
1722+
private registry.
1723+
items:
1724+
description: |-
1725+
LocalObjectReference contains enough information to let you locate the
1726+
referenced object inside the same namespace.
1727+
properties:
1728+
name:
1729+
default: ""
1730+
description: |-
1731+
Name of the referent.
1732+
This field is effectively required, but due to backwards compatibility is
1733+
allowed to be empty. Instances of this type with an empty value here are
1734+
almost certainly wrong.
1735+
More info: https://kubernetes.io/docs/concepts/overview/working-with-objects/names/#names
1736+
type: string
1737+
type: object
1738+
x-kubernetes-map-type: atomic
1739+
type: array
17081740
nodeSelector:
17091741
additionalProperties:
17101742
type: string
17111743
description: NodeSelector is a selector which must be true for the
17121744
pod to fit on a node
17131745
type: object
1746+
processorImage:
1747+
description: |-
1748+
ProcessorImage is the full container image for the dataflow processor (e.g. ghcr.io/org/dataflow:v1.0.0).
1749+
If not set, the processor runs with the same image as the controller (or ProcessorVersion is used if set).
1750+
type: string
1751+
processorVersion:
1752+
description: |-
1753+
ProcessorVersion is the image tag for the processor when using the default image repository (same as controller).
1754+
Ignored if ProcessorImage is set. Example: "v1.2.3".
1755+
type: string
17141756
resources:
17151757
description: Resources defines the resource requirements for the processor
17161758
pod
@@ -2190,7 +2232,8 @@ spec:
21902232
format: int32
21912233
type: integer
21922234
batchSize:
2193-
description: BatchSize for batch inserts
2235+
default: 1
2236+
description: 'BatchSize for batch inserts (default: 1)'
21942237
format: int32
21952238
type: integer
21962239
conflictKey:
@@ -2219,6 +2262,16 @@ spec:
22192262
- key
22202263
- name
22212264
type: object
2265+
rawMode:
2266+
description: |-
2267+
RawMode when true, expects messages in format {"value": <data>, "_metadata": {...}}. Table is created with value JSONB and _metadata JSONB columns.
2268+
When false, table structure is inferred from the first message (replicates source structure).
2269+
type: boolean
2270+
softDeleteColumn:
2271+
description: SoftDeleteColumn specifies column for soft delete
2272+
(e.g. "deleted_at"). If set, DELETE operations will UPDATE
2273+
this column instead of physical delete.
2274+
type: string
22222275
table:
22232276
description: Table to write to
22242277
type: string
@@ -3138,6 +3191,16 @@ spec:
31383191
postgresql:
31393192
description: PostgreSQL source configuration
31403193
properties:
3194+
autoCreateTable:
3195+
description: AutoCreateTable creates the table if it doesn't
3196+
exist before reading
3197+
type: boolean
3198+
changeTrackingColumn:
3199+
default: updated_at
3200+
description: |-
3201+
ChangeTrackingColumn is the column used to track changes (default: updated_at).
3202+
Not used when Query is specified.
3203+
type: string
31413204
connectionString:
31423205
description: ConnectionString for PostgreSQL database
31433206
type: string
@@ -3171,6 +3234,11 @@ spec:
31713234
description: 'RawMode when true, wraps each row as JSON: {"value":
31723235
<row data>, "_metadata": {table, id}}'
31733236
type: boolean
3237+
readBatchSize:
3238+
description: ReadBatchSize limits rows per poll to reduce
3239+
DB load (0 = no limit)
3240+
format: int32
3241+
type: integer
31743242
table:
31753243
description: Table to read from
31763244
type: string
@@ -3510,30 +3578,6 @@ spec:
35103578
type: string
35113579
type: object
35123580
type: array
3513-
processorImage:
3514-
description: ProcessorImage is the full container image for the dataflow
3515-
processor (e.g. ghcr.io/org/dataflow:v1.0.0). If not set, the processor
3516-
runs with the same image as the controller (or ProcessorVersion is
3517-
used if set).
3518-
type: string
3519-
processorVersion:
3520-
description: ProcessorVersion is the image tag for the processor when
3521-
using the default image repository (same as controller). Ignored if
3522-
ProcessorImage is set. Example "v1.2.3".
3523-
type: string
3524-
imagePullSecrets:
3525-
description: ImagePullSecrets is a list of references to secrets in
3526-
the same namespace to use for pulling the processor image from a
3527-
private registry.
3528-
items:
3529-
description: LocalObjectReference contains enough information to
3530-
let you locate the referenced object inside the same namespace.
3531-
properties:
3532-
name:
3533-
description: Name of the referent.
3534-
type: string
3535-
type: object
3536-
type: array
35373581
transformations:
35383582
description: Transformations is a list of transformations to apply
35393583
to messages
@@ -4065,7 +4109,9 @@ spec:
40654109
format: int32
40664110
type: integer
40674111
batchSize:
4068-
description: BatchSize for batch inserts
4112+
default: 1
4113+
description: 'BatchSize for batch inserts
4114+
(default: 1)'
40694115
format: int32
40704116
type: integer
40714117
conflictKey:
@@ -4096,6 +4142,17 @@ spec:
40964142
- key
40974143
- name
40984144
type: object
4145+
rawMode:
4146+
description: |-
4147+
RawMode when true, expects messages in format {"value": <data>, "_metadata": {...}}. Table is created with value JSONB and _metadata JSONB columns.
4148+
When false, table structure is inferred from the first message (replicates source structure).
4149+
type: boolean
4150+
softDeleteColumn:
4151+
description: SoftDeleteColumn specifies column
4152+
for soft delete (e.g. "deleted_at"). If
4153+
set, DELETE operations will UPDATE this
4154+
column instead of physical delete.
4155+
type: string
40994156
table:
41004157
description: Table to write to
41014158
type: string

‎config/samples/pg-to-pg-test.yaml‎

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
2+
apiVersion: dataflow.dataflow.io/v1
3+
kind: DataFlow
4+
metadata:
5+
name: pg-to-pg-test
6+
spec:
7+
source:
8+
type: postgresql
9+
postgresql:
10+
connectionString: "postgres://dataflow:dataflow@localhost:5432/dataflow?sslmode=disable"
11+
table: public.products
12+
pollInterval: 5
13+
sink:
14+
type: postgresql
15+
postgresql:
16+
connectionString: "postgres://dataflow:dataflow@localhost:5432/dataflow?sslmode=disable"
17+
table: public.products_clone
18+
autoCreateTable: true
19+
upsertMode: true
20+
batchSize: 10
21+
batchFlushIntervalSeconds: 2

‎config/samples/pg-to-pg-test2.yaml‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
2+
apiVersion: dataflow.dataflow.io/v1
3+
kind: DataFlow
4+
metadata:
5+
name: pg-to-pg-test
6+
spec:
7+
source:
8+
type: postgresql
9+
postgresql:
10+
connectionString: "postgres://dataflow:dataflow@localhost:5432/dataflow?sslmode=disable"
11+
table: public.products
12+
pollInterval: 5
13+
sink:
14+
type: postgresql
15+
postgresql:
16+
connectionString: "postgres://dataflow:dataflow@localhost:5432/dataflow?sslmode=disable"
17+
table: public.products_raw_clone
18+
autoCreateTable: true
19+
rawMode: true
20+
upsertMode: true
21+
batchSize: 10
22+
batchFlushIntervalSeconds: 2

0 commit comments

Comments
 (0)