apache/seatunnel · error · NatsJetStreamConnectorException

ILLEGAL_ARGUMENT

ILLEGAL_ARGUMENT

Error message

Invalid NATS JetStream sink option `format`: must be one of: json, native

What it means

NatsJetStreamSinkValidator.validate requires the `format` option to resolve to a known NatsJetStreamMessageFormat (json or native). A missing or null format fails validation with ILLEGAL_ARGUMENT before the sink starts.

Source

Thrown at seatunnel-connectors-v2/connector-nats-jetstream/src/main/java/org/apache/seatunnel/connectors/seatunnel/natsjetstream/sink/NatsJetStreamSinkValidator.java:54

final class NatsJetStreamSinkValidator {

    private static final Set<String> SUPPORTED_NATIVE_MAPPING_KEYS =
            new HashSet<>(
                    Arrays.asList(
                            NatsJetStreamSinkOptions.NATIVE_MAPPING_ID,
                            NatsJetStreamSinkOptions.NATIVE_MAPPING_SUBJECT,
                            NatsJetStreamSinkOptions.NATIVE_MAPPING_HEADERS,
                            NatsJetStreamSinkOptions.NATIVE_MAPPING_DATA));

    private NatsJetStreamSinkValidator() {}

    static void validate(ReadonlyConfig pluginConfig, CatalogTable catalogTable) {
        requireNonBlank(pluginConfig, NatsJetStreamSinkOptions.URL);
        validateAuthentication(pluginConfig);
        NatsJetStreamMessageFormat format = pluginConfig.get(NatsJetStreamSinkOptions.FORMAT);
        if (format == null) {
            throw invalidOption("format", "must be one of: json, native");
        }
        if (format == NatsJetStreamMessageFormat.JSON) {
            requireNonBlank(pluginConfig, NatsJetStreamSinkOptions.SUBJECT);
            return;
        }
        validateNativeMapping(pluginConfig, catalogTable);
    }

    private static void validateAuthentication(ReadonlyConfig pluginConfig) {
        Optional<String> username = pluginConfig.getOptional(NatsJetStreamSinkOptions.USERNAME);
        Optional<String> password = pluginConfig.getOptional(NatsJetStreamSinkOptions.PASSWORD);
        Optional<String> token = pluginConfig.getOptional(NatsJetStreamSinkOptions.TOKEN);

        boolean hasUsername = username.map(NatsJetStreamSinkValidator::isNotBlank).orElse(false);
        boolean hasPassword = password.map(NatsJetStreamSinkValidator::isNotBlank).orElse(false);
        boolean hasToken = token.map(NatsJetStreamSinkValidator::isNotBlank).orElse(false);

        if (hasUsername != hasPassword) {

View on GitHub (pinned to cf67b549a7)

Solutions

  1. Set format = json for plain JSON messages
  2. Set format = native to publish column-mapped native payloads (also requires native_format_fields)
  3. Check spelling/case of the format value against NatsJetStreamSinkOptions.FORMAT allowed values

Example fix

// before
sink {
  NatsJetStream {
    url = "nats://localhost:4222"
  }
}
// after
sink {
  NatsJetStream {
    url = "nats://localhost:4222"
    format = json
    subject = "events"
  }
}
Defensive patterns

Strategy: validation

Validate before calling

Set<String> allowed = Set.of("json", "native");
if (!allowed.contains(config.getString("format"))) {
    throw new IllegalArgumentException("format must be json or native");
}

Prevention

When it happens

Trigger: format option omitted, misspelled (e.g. 'JSON', 'text'), or set to a value not supported by the NATS JetStream sink in this SeaTunnel version.

Common situations: Copy-pasted config from another sink (e.g. Kafka) that uses different format names; version upgrade where format keys changed; typo like 'naive' for 'native'.

Understand the failure class

Background: Invalid enum value errors: "Unknown type", "Invalid scope", "must be one of" — when a string is not on the library's allowed list — this error's family across 23 libraries.

Related errors


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