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
- Inspect the wrapped cause (getCause()) for the real RocketMQ client error
- Ensure producer shutdown completes/flushes before close and that clients aren't closed twice in multi-writer setups
- 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
- Make close idempotent and guard with a closed flag
- Never let close() mask the original job exception; log-and-warn instead
- Test shutdown paths (multi-writer shared clients) to avoid double-close
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
- WRITER_CLOSE_FAILED
- CLOSE_FAILED
- CLOSE_CONNECTION_FAILED
- Failed to close HugeGraph sink writer
- Failed to close %s via JDBC.
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/46515941878b7b15.
Report an issue: GitHub.