Sobes.tech
Middle+

/* Nous devons transférer des données d'une source à un consommateur. La source fournit les données en petits lots (~dix enregistrements), tandis que le consommateur fonctionne plus efficacement avec de gros lots (~mille enregistrements). Un exemple réel est le transfert de données depuis des files d'attente de type Kafka vers une base de données Clickhouse. Source : - Pratiquement infinie. - La source ne renvoie jamais plus de MaxItems enregistrements en un seul appel à Next. - Dans le cadre d'une "session" (un appel à la fonction Pipe), la source renvoie de nouvelles données à chaque appel à Next. - Cependant, après un redémarrage, la source recommence à partir de la position "confirmée" précédente, définie par cookie. Par conséquent, *chaque* valeur de cookie renvoyée par Next, après avoir sauvegardé les données dans le récepteur, doit être confirmée avec un appel à Commit, dans le même ordre dans lequel elles ont été renvoyées par Next. Récepteur : - Ne peut traiter plus de MaxItems à la fois. Il faut implémenter la fonction func Pipe(p Producer, c Consumer) error qui lit les données de la source, les regroupe dans un buffer de taille maximale MaxItems et les sauvegarde dans le récepteur, après quoi elle confirme la progression dans la source. */ const MaxItems = 9999 type Producer interface { // Next retourne : // - un lot d'éléments à traiter // - cookie à confirmer une fois le traitement terminé // - erreur Next() (items []any, cookie int, err error) // Commit est utilisé pour marquer le lot de données comme traité Commit(cookie int) error } type Consumer interface { Process(items []any) error } func Pipe(p Producer, c Consumer) error { // TODO }