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.