apache/seatunnel · error · SnmpConnectorException

INVALID_ROW

INVALID_ROW

Error message

Input row arity ${arity} does not match the configured schema arity ${rowArity}

What it means

The SNMP sink validates each incoming SeaTunnelRow against the schema arity declared in the sink configuration before performing any network I/O. If the row's field count differs from the configured row type, conversion is rejected with INVALID_ROW to prevent index misalignment when reading the oid/value fields.

Solutions

  1. Align the upstream schema so the row arity matches the sink's configured SeaTunnelRowType
  2. Check the source/transform column list and remove or add columns to match the SNMP sink table definition
  3. Rebuild the pipeline after schema edits — old checkpoints may replay rows with the old arity, so restart with a fresh state if needed
  4. Wrap the sink write path in a catalog/schema consistency check at job submit time

Example fix

// before: source SELECT with 3 cols, sink schema with 2
SELECT host, oid, value FROM events
// after
SELECT oid, value FROM events  -- matches sink schema arity 2
Defensive patterns

Strategy: validation

Validate before calling

// Before the sink write, assert arity:
if (row.getArity() != sinkRowType.getTotalFields()) {
    throw new IllegalArgumentException("row arity " + row.getArity()
        + " != schema arity " + sinkRowType.getTotalFields());
}

Try / catch

try {
    sink.write(row);
} catch (SeaTunnelException e) {
    if (e.getMessage().contains("does not match the configured schema arity")) {
        LOG.error("Upstream schema drift detected for SNMP sink", e);
    }
    throw e;
}

Prevention

When it happens

Trigger: SnmpSinkRowConverter.convert(row) is called (from the sink's write path) with a SeaTunnelRow whose getArity() != rowArity — e.g. an upstream transform changed the number of columns, or the sink table schema was edited without updating the upstream source.

Common situations: Adding/removing a column in the source SQL or transform upstream of the SNMP sink; two sinks sharing one dataset with different schemas; row produced by a custom Source with a stale SeaTunnelRowType.

Understand the failure class

Background: Schema validation failed / invalid input schema: payload rejected because its shape doesn't match the expected schema — this error's family across 28 libraries.

Related errors


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

Appendix: source

Thrown at seatunnel-connectors-v2/connector-snmp/src/main/java/org/apache/seatunnel/connectors/seatunnel/snmp/sink/SnmpSinkRowConverter.java:80

    private static final BigInteger TICKS_PER_MINUTE = BigInteger.valueOf(6_000L);
    private static final BigInteger TICKS_PER_SECOND = BigInteger.valueOf(100L);

    private final int rowArity;
    private final int oidIndex;
    private final int valueIndex;
    private final int valueTypeIndex;

    SnmpSinkRowConverter(SnmpSinkConfig config, SeaTunnelRowType rowType) {
        this.rowArity = rowType.getTotalFields();
        this.oidIndex = requireStringField(rowType, config.getOidField(), "oid_field");
        this.valueIndex = requireStringField(rowType, config.getValueField(), "value_field");
        this.valueTypeIndex =
                requireStringField(rowType, config.getValueTypeField(), "value_type_field");
    }

    SnmpSetRequest convert(SeaTunnelRow row) {
        if (row.getArity() != rowArity) {
            throw invalidRow(
                    "Input row arity "
                            + row.getArity()
                            + " does not match the configured schema arity "
                            + rowArity);
        }

        String oid = requireNonBlankRowValue(row, oidIndex, "OID");
        String value = requireNonNullRowValue(row, valueIndex, "value");
        String valueType = requireNonBlankRowValue(row, valueTypeIndex, "value type");
        return new SnmpSetRequest(parseOid(oid), parseVariable(valueType, value));
    }

    private static int requireStringField(
            SeaTunnelRowType rowType, String fieldName, String optionName) {
        int index = rowType.indexOf(fieldName, false);
        if (index < 0) {
            throw invalidConfig(
                    "Option `"

View on GitHub (pinned to cf67b549a7)