Skip to content

Commit 0b81110

Browse files
committed
Use new Upsert and Delete structs
1 parent 591a806 commit 0b81110

File tree

1 file changed

+6
-6
lines changed

1 file changed

+6
-6
lines changed

pkg/sync/metrics.go

+6-6
Original file line numberDiff line numberDiff line change
@@ -118,11 +118,11 @@ func (ms *MetricSync) Run(ctx context.Context) error {
118118
})
119119

120120
g.Go(func() error {
121-
return ms.db.UpsertStreamed(ctx, upsertPodMetrics, database.WithStatement(ms.podMetricUpsertStmt(), 5))
121+
return database.NewUpsert(ms.db).WithStatement(ms.podMetricUpsertStmt(), 5).Stream(ctx, upsertPodMetrics)
122122
})
123123

124124
g.Go(func() error {
125-
return ms.db.UpsertStreamed(ctx, upsertContainerMetrics, database.WithStatement(ms.containerMetricUpsertStmt(), 6))
125+
return database.NewUpsert(ms.db).WithStatement(ms.containerMetricUpsertStmt(), 6).Stream(ctx, upsertContainerMetrics)
126126
})
127127

128128
return g.Wait()
@@ -157,11 +157,11 @@ func (ms *MetricSync) Clean(ctx context.Context, deleteChannel <-chan contracts.
157157
})
158158

159159
g.Go(func() error {
160-
return ms.db.DeleteStreamed(ctx, &schema.PodMetric{}, deletesPod, database.ByColumn("reference_id"))
160+
return database.NewDelete(ms.db).ByColumn("reference_id").Stream(ctx, &schema.PodMetric{}, deletesPod)
161161
})
162162

163163
g.Go(func() error {
164-
return ms.db.DeleteStreamed(ctx, &schema.ContainerMetric{}, deletesContainer, database.ByColumn("pod_reference_id"))
164+
return database.NewDelete(ms.db).ByColumn("pod_reference_id").Stream(ctx, &schema.ContainerMetric{}, deletesContainer)
165165
})
166166

167167
return g.Wait()
@@ -229,7 +229,7 @@ func (nms *NodeMetricSync) Run(ctx context.Context) error {
229229
})
230230

231231
g.Go(func() error {
232-
return nms.db.UpsertStreamed(ctx, upsertNodeMetrics, database.WithStatement(nms.nodeMetricUpsertStmt(), 5))
232+
return database.NewUpsert(nms.db).WithStatement(nms.nodeMetricUpsertStmt(), 5).Stream(ctx, upsertNodeMetrics)
233233
})
234234

235235
return g.Wait()
@@ -261,7 +261,7 @@ func (nms *NodeMetricSync) Clean(ctx context.Context, deleteChannel <-chan contr
261261
})
262262

263263
g.Go(func() error {
264-
return nms.db.DeleteStreamed(ctx, &schema.NodeMetric{}, deletes, database.ByColumn("node_id"))
264+
return database.NewDelete(nms.db).ByColumn("node_id").Stream(ctx, &schema.NodeMetric{}, deletes)
265265
})
266266

267267
return g.Wait()

0 commit comments

Comments
 (0)