Incident critique d'erreur CRC dans un environnement Kafka

Problème initial

En production, le système d'un client a commencé à recevoir de manière intermittente des exceptions liées au CRC dans Apache Kafka. L'erreur suivante apparaît :


Record batch for partition skywalking-traces-0 at offset 292107075 is invalid, cause: Record is corrupt (stored crc = 1016021496, compute crc = 1981017560)

Cette erreur est imprévisible : elle peut survenir plusieurs fois par heure ou pas du tout pendant une demi-journée. Elle provoque le blocage des consommateurs Kafka, les empêchant de traiter les messages ultérieurs. Un redémarrage manuel est nécessaire pour reprendre le traitement, ce qui entraîne une indisponibilité en production.

Pour comprendre cette erreur, rappelons que Kafka stocke un CRC (Cyclic Redundancy Check) dans l'en-tête de chaque lot de messages. Lorsqu'un consommateur reçoit un message, il calcule à nouveau le CRC à partir du contenu et le compare à celui stocké. Une divergence indique une corruption des données.

Analyse et localisation

Comparaison avec un autre cas

Une autre application (client A) a rencontré la même erreur de CRC, mais avec un comportement différent. Pour le client A, l'erreur est reproductible de manière stable : une fois apparue, elle persiste même après des redémarrages multiples des consommateurs. Le problème a été tracé à un secteur défectueux sur le disque dur, résolu après remplacement.

En revanche, pour le client B, l'erreur est éphémère : elle disparaît après un redémarrage du consommateur ou lors du basculement vers un autre consommateur. Elle ne peut pas être reproduite de manière fiable.

Exclusion des bogues Kafka

En tant que contributeur actif à la communauté Kafka et avec une expérience sur le cloud public d'Alibaba, je n'ai pas rencontré cette erreur auparavant. Les problèmes de CRC étaient courants dans les versions antérieures à 0.10, mais le client B utilise la version 2.8.2, une version stable. Le code de calcul du CRC dans Kafka est succinct (quelques dizaines de lignes de Java), ce qui rend un bogue Kafka peu probable. L'effort a donc été redirigé vers d'autres causes.

Caractéristiques de l'erreur

L'analyse a révélé les caractéristiques suivantes :

  • L'erreur est intermittente, probabiliste et imprévisible.
  • Elle ne peut pas être reproduite : les messages défectueux semblent sains lors d'une nouvelle tentative de consommation.
  • Seul un environnement parmi huit est affecté, malgré des configurations matérielles similaires.
  • Aucune corrélation n'a été trouvée avec des brokers ou des hôtes spécifiques.
  • Les journaux système n'ont pas révélé d'erreurs de mémoire.
  • Le système d'exploitation personnalisé utilise la pile TCP native de Linux, sans modifications.
  • Les simulations de perte ou d'erreur réseau en environnement de test n'ont pas réduit le problème.

Les efforts se sont concentrés sur le réseau et le stockage, des composants généralement fiables.

Investigation réseau

Capture de trafic

Pour examiner le trafic réseau, des captures TCP ont été déployées sur les trois brokers (10.0.0.70, 10.0.0.71, 10.0.0.72) et un consommateur (10.0.0.18). Les commandes suivantes ont été exécutées pour limiter la taille des fichiers de capture :


# Sur le consommateur 10.0.0.18
sudo nohup tcpdump -i vethfe2axxxx -nve host 10.0.0.70 or host 10.0.0.71 or host 10.0.0.72 -w broker_18.pcap -C 500 -W 5 -Z ccadmin &

# Sur les brokers (exemple pour 10.0.0.70)
sudo nohup tcpdump -i eth0 -nve host 10.0.0.18 -w consumer_18.pcap -C 500 -W 5 -Z ccadmin &

Les options -C 500 et -W 5 divisent les captures en fichiers de 500 Mo et conservent les cinq plus récents.

Détection d'anomalies

Après deux heures d'attente, une erreur a été capturée. L'analyse avec Wireshark a montré une divergence de 15 octets dans un paquet de 2712 octets. Toutefois, Wireshark ne permet pas de valider le contenu Kafka en profondeur. Des comparaisons avec d'autres messages normaux ont confirmé l'intégrité des données avant transmission.

Décodage du protocole Kafka

Pour vérifier l'intégrité, le protocole Kafka a été décodé manuellement à partir des captures binaires. Le code suivant analyse une requête Fetch :


import org.apache.kafka.common.requests.FetchRequest;
import org.apache.kafka.common.requests.RequestHeader;
import javax.xml.bind.DatatypeConverter;
import java.nio.ByteBuffer;

public class KafkaRequestAnalyzer {
    public static void main(String[] args) {
        // Données hexadécimales extraites de la capture réseau
        String hexPayload = "0e 1f 48 2d 7e 32 06 82 25 e9 a3 d9 08 00 45 00 " +
                            "00 ab 3f b3 40 00 40 06 35 ed 62 02 00 55 62 02 " +
                            "00 54 eb 2e 23 85 68 57 32 37 b5 08 3f b2 80 18 " +
                            "7d 2c c5 4a 00 00 01 01 08 0a b7 d1 e6 5c 26 30 " +
                            "b9 11 00 00 00 73 00 01 00 0c 11 85 6f bf 00 15 " +
                            "62 72 6f 6b 65 72 2d 31 30 30 32 2d 66 65 74 63 " +
                            "68 65 72 2d 30 00 00 00 03 ea 00 00 01 f4 00 00 " +
                            "00 01 00 a0 00 00 00 72 52 11 e0 11 85 6f bf 02 " +
                            "13 5f 5f 63 6f 6e 73 75 6d 65 72 5f 6f 66 66 73 " +
                            "65 74 73 02 00 00 00 15 00 00 00 03 00 00 00 00 " +
                            "13 36 67 3a 00 00 00 03 00 00 00 00 00 00 00 00 " +
                            "00 10 00 00 00 00 01 01 00";

        // Nettoyage et conversion en tableau d'octets
        String cleanedHex = hexPayload.replaceAll("\\s+", "");
        byte[] binaryData = DatatypeConverter.parseHexBinary(cleanedHex);
        ByteBuffer buffer = ByteBuffer.wrap(binaryData);

        // Ignorer l'en-tête Ethernet et IP (66 octets) + longueur (4 octets)
        buffer.position(70);

        // Décodage de l'en-tête de requête Kafka
        RequestHeader reqHeader = RequestHeader.parse(buffer);
        System.out.println("En-tête de requête : " + reqHeader);

        // Décodage de la requête Fetch
        FetchRequest fetchReq = FetchRequest.parse(buffer, (short) 12);
        System.out.println("Requête Fetch : " + fetchReq.toStruct());
    }
}

La sortie a confirmé que le contenu de la requête était identique avant et après transmission, éliminant un problème de requête.

Une analyse similaire a été effectuée pour la réponse, révélant que le CRC calculé dynamiquement à partir des données reçues ne correspondait pas à celui stocké, alors qu'il correspondait avant envoi. Cela suggérait une altération en transit.

Investigation disque

Les données sur le disque ont été comparées aux données envoyées. Aucune divergence n'a été trouvée, éliminant le disque comme source du problème.

Conclusion

La divergence dans les données TCP avant et après transmission, uniquement pour les messages défectueux, a permis de localiser le problème au niveau du réseau. L'équipe réseau a confirmé que ce comportement était anormal et a lancé une investigation approfondie.

Épilogue

L'enquête a révélé que la carte réseau d'un fabricant spécifique modifiait de manière probabiliste les paquets lors de l'opération TCP Segmentation Offload (TSO). Le checksum TCP était préservé car les modifications étaient cohérentes avec la logique de calcul, permettant aux paquets corrompus de passer la validation. Le fabricant a reconnu le défaut matériel et travaille sur une correction.

Étiquettes: Apache Kafka CRC Réseau TCP Débogage réseau Protocole Kafka

Publié le 28 juillet à 04h30