/*
Πρέπει να μεταφέρουμε δεδομένα από μια πηγή σε έναν καταναλωτή.
Η πηγή παρέχει δεδομένα σε μικρές παρτίδες (~δέκα εγγραφές), ενώ ο καταναλωτής λειτουργεί πιο αποδοτικά με μεγάλες παρτίδες (~χίλιες εγγραφές).
Ένα πραγματικό παράδειγμα είναι η μεταφορά δεδομένων από ουρές τύπου Kafka σε μια βάση δεδομένων Clickhouse.
Πηγή:
- Πρακτικά άπειρη.
- Η πηγή ποτέ δεν επιστρέφει περισσότερα από MaxItems εγγραφές σε μια κλήση Next.
- Σε μια "συνεδρία" (μία κλήση της συνάρτησης Pipe), η πηγή επιστρέφει νέα δεδομένα σε κάθε κλήση Next.
- Ωστόσο, μετά από επανεκκίνηση, η πηγή ξεκινά από την προηγούμενη "επιβεβαιωμένη" θέση, ορισμένη από το cookie.
Επομένως, *κάθε* τιμή cookie που επιστρέφει το Next, μετά την αποθήκευση των δεδομένων στον δέκτη,
πρέπει να επιβεβαιωθεί με μια κλήση Commit, με την ίδια σειρά με αυτή που επιστράφηκε από το Next.
Δέκτης:
- Δεν μπορεί να επεξεργαστεί περισσότερα από MaxItems ταυτόχρονα.
Απαιτείται η υλοποίηση της συνάρτησης func Pipe(p Producer, c Consumer) error
που διαβάζει δεδομένα από την πηγή, τα ομαδοποιεί σε ένα buffer μεγέθους όχι μεγαλύτερο από MaxItems και τα αποθηκεύει στον δέκτη,
μετά επιβεβαιώνει την πρόοδο στην πηγή.
*/
const MaxItems = 9999
type Producer interface {
// Next επιστρέφει:
// - ένα batch από items προς επεξεργασία
// - cookie για επιβεβαίωση όταν ολοκληρωθεί η επεξεργασία
// - σφάλμα
Next() (items []any, cookie int, err error)
// Commit χρησιμοποιείται για να σηματοδοτήσει το batch δεδομένων ως επεξεργασμένο
Commit(cookie int) error
}
type Consumer interface {
Process(items []any) error
}
func Pipe(p Producer, c Consumer) error {
// TODO
}