apache/seatunnel · error · RocketMqConnectorException

WRITER_OPERATION_FAILED

WRITER_OPERATION_FAILED

Error message

Close RocketMq sink writer error

What it means

RocketMqSinkWriter.close closes the underlying RocketMQ producer sender. If sender.close() throws any Exception, it is rewrapped as RocketMqConnectorException with WRITER_OPERATION_FAILED and the message 'Close RocketMq sink writer error', preserving the original as the cause.

Source

Thrown at seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/sink/RocketMqSinkWriter.java:69

        }
        // Set `rocketmq.client.logUseSlf4j` to `true` to avoid create many
        // `AsyncAppender-Dispatcher-Thread`
        System.setProperty("rocketmq.client.logUseSlf4j", "true");
    }

    @Override
    public void write(SeaTunnelRow element) throws IOException {
        Message message = seaTunnelRowSerializer.serializeRow(element);
        rocketMqProducerSender.send(message);
    }

    @Override
    public void close() throws IOException {
        if (this.rocketMqProducerSender != null) {
            try {
                this.rocketMqProducerSender.close();
            } catch (Exception e) {
                throw new RocketMqConnectorException(
                        CommonErrorCodeDeprecated.WRITER_OPERATION_FAILED,
                        "Close RocketMq sink writer error",
                        e);
            }
        }
    }

    private SeaTunnelRowSerializer<byte[], byte[]> getSerializer(
            SeaTunnelRowType seaTunnelRowType) {
        return new DefaultSeaTunnelRowSerializer(
                producerMetadata.getTopic(),
                producerMetadata.getTag(),
                getPartitionKeyFields(seaTunnelRowType),
                seaTunnelRowType,
                producerMetadata.getFormat(),
                producerMetadata.getFieldDelimiter());
    }

View on GitHub (pinned to cf67b549a7)

Solutions

  1. Inspect the wrapped cause (getCause()) for the real RocketMQ client error
  2. Ensure producer shutdown completes/flushes before close and that clients aren't closed twice in multi-writer setups
  3. Verify broker connectivity and producer shutdownTimeout settings; log-and-suppress in close if failing close shouldn't mask the original job error

Example fix

// before
public void close() throws IOException {
    this.rocketMqProducerSender.close();
}
// after
public void close() throws IOException {
    try {
        this.rocketMqProducerSender.close();
    } catch (Exception e) {
        log.warn("Error closing RocketMQ producer", e);
    }
}
Defensive patterns

Strategy: try-catch

Validate before calling

if (rocketMqProducerSender != null && !rocketMqProducerSender.isClosed()) rocketMqProducerSender.close();

Try / catch

try { writer.close(); } catch (RocketMqConnectorException e) { log.warn("RocketMQ producer close failed", e.getCause()); }

Prevention

When it happens

Trigger: Writer cleanup (job shutdown, checkpoint teardown, or test code closing writers) when rocketMqProducerSender.close() fails — e.g. in-flight send shutdown errors or producer client cleanup exceptions.

Common situations: Producer already half-closed due to an earlier failure; network problems during producer shutdown; broker unreachable while flushing on close; shared client closed twice across multi-table writers.

Related errors


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