apache/seatunnel · error · EdgeSocketConnectorException

DECODE_FAILED

DECODE_FAILED

Error message

Decode ingress packet failed for batchId={}

What it means

WARN log in handleBatchRecord for any EdgeSocketConnectorException that is NOT a decryption error: the ingress packet could not be decoded (framing/serialization failure). The batch is rejected with the DECODE_FAILED response code and the sender is expected to fix and resend or drop the batch.

Solutions

  1. Check the logged stack trace to identify the decode stage that failed (header parse vs payload deserialize)
  2. Align the sender's serialization format/protocol version with the connector's expected packet format
  3. Re-snapshot the batch at the device: truncated packets usually indicate a sender-side write that was cut short
  4. Upgrade or roll back connector/client versions so both sides speak the same wire protocol
  5. Send a known-good sample batch from a test client to validate the path end to end

Example fix

// before: sender uses incompatible format
client.send(payloadBytes, Format.JSON);
// after: match connector's expected format
client.send(packetCodec.encode(batch), Format.CONNECTOR_V1);
Defensive patterns

Strategy: retry

Validate before calling

// sender-side: validate packet framing before send
if (payload == null || payload.length == 0) throw new IllegalArgumentException("empty batch payload");

Try / catch

// on DECODE_FAILED, do not blind-retry the same bytes; re-encode once
if ("DECODE_FAILED".equals(responseCode)) {
    byte[] reEncoded = packetCodec.encode(batch);
    client.sendBatch(batchId, reEncoded);
}

Prevention

When it happens

Trigger: handleBatchRecord decode step throws EdgeSocketConnectorException with a non-decryption cause — malformed wire format, unsupported serialization type, truncated packet, protocol version mismatch.

Common situations: Edge client upgraded to a newer wire protocol the connector doesn't understand; partial reads / corrupted TCP stream framing; sender writing plain JSON where connector expects its binary format (or vice versa); buggy device firmware producing malformed batches.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/477d7390e6cc6a55. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-edge-socket/src/main/java/org/apache/seatunnel/connectors/seatunnel/edgesocket/source/EdgeSocketSourceReader.java:208

            queueFullCount++;
            if (queueFullCount == 1 || queueFullCount % 100 == 0) {
                log.warn(
                        "Ingress queue physically full, returning RETRY "
                                + "(capacity={}, rejectCount={}, batchId={})",
                        config.getLocalQueueCapacity(),
                        queueFullCount,
                        batchId);
            }
            return EdgeSocketResponseCode.RETRY.getCode();
        } catch (EdgeSocketConnectorException connectorException) {
            if (isDecryptionError(connectorException)) {
                log.warn(
                        "Decryption failed for batchId={}, check secret_key configuration",
                        batchId,
                        connectorException);
                return EdgeSocketResponseCode.DECRYPT_FAILED.getCode();
            }
            log.warn("Decode ingress packet failed for batchId={}", batchId, connectorException);
            return EdgeSocketResponseCode.DECODE_FAILED.getCode();
        } catch (Exception decodeException) {
            log.warn("Decode or enqueue ingress packet failed", decodeException);
            return EdgeSocketResponseCode.DECODE_FAILED.getCode();
        }
    }

    @Override
    public String handleCommitRequest(long batchId) {
        synchronized (stateLock) {
            return sourceState.resolveCommitResponse(batchId);
        }
    }

    private boolean isDecryptionError(EdgeSocketConnectorException exception) {
        EdgeSocketConnectorErrorCode errorCode =
                (EdgeSocketConnectorErrorCode) exception.getSeaTunnelErrorCode();
        return errorCode == EdgeSocketConnectorErrorCode.PACKET_AES_KEY_MISSING

View on GitHub (pinned to cf67b549a7)