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
- Align the upstream schema so the row arity matches the sink's configured SeaTunnelRowType
- Check the source/transform column list and remove or add columns to match the SNMP sink table definition
- Rebuild the pipeline after schema edits — old checkpoints may replay rows with the old arity, so restart with a fresh state if needed
- 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
- Keep sink table schemas and upstream source/transform column lists in one reviewed place
- Re-run schema validation after every pipeline edit
- Avoid sharing one dataset between sinks with different schemas
- Restart from clean state after schema changes to avoid replayed old-arity rows
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)