risingwavelabs/risingwave · error · anyhow

sink validation failed

Error message

sink validation failed: {}

What it means

After RisingWave sends a validation request to the remote (JVM) sink via JNI, the remote side returns a ValidateSinkResponse containing an error message. RisingWave surfaces that remote-side failure verbatim as 'sink validation failed: <message>'.

Solutions

  1. Read the inner error_message appended to this error - it comes from the remote connector and names the actual problem.
  2. Fix the sink's connection properties (endpoint, credentials, format/encode options).
  3. Verify the downstream system is reachable from the RisingWave node.
  4. Re-run validation after correcting the CREATE SINK options.
Defensive patterns

Strategy: try-catch

Try / catch

match sink_validate_result {
    Err(e) if e.to_string().contains("sink validation failed") => {
        // surface the inner remote message, fix connector properties/endpoint
        log::error("remote sink rejected validation: {}", e);
    }
    Err(e) => return Err(e),
    Ok(()) => { /* proceed */ }
}

Prevention

When it happens

Trigger: Calling sink validate where the remote connector (running in the embedded JVM) rejects the sink config/schema, returning error set in ValidateSinkResponse.error.error_message.

Common situations: Wrong connection settings for the downstream system (bad ES URL, JDBC credentials, Kafka broker address); remote connector rejects the schema during its own validation pass; the downstream service is unreachable from the JVM.

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 risingwavelabs/risingwave@6469eb736d (2026-09-11). Data as JSON: /api/errors/335241aed80fb014. Report an issue: GitHub.

Appendix: source

Thrown at src/connector/src/sink/remote.rs:265

                sink_param: Some(sink_param),
            };
            let validate_sink_request_bytes =
                env.byte_array_from_slice(&Message::encode_to_vec(&validate_sink_request))?;

            let validate_sink_response_bytes = call_static_method!(
                env,
                {com.risingwave.connector.JniSinkValidationHandler},
                {byte[] validate(byte[] validateSourceRequestBytes)},
                &validate_sink_request_bytes
            )?;

            let validate_sink_response: ValidateSinkResponse = Message::decode(
                risingwave_jni_core::to_guarded_slice(&validate_sink_response_bytes, env)?.deref(),
            )?;

            validate_sink_response.error.map_or_else(
                || Ok(()), // If there is no error message, return Ok here.
                |err| bail!("sink validation failed: {}", err.error_message),
            )
        })
    })
    .await
    .context("join handle returned an error")??;

    Ok(())
}

pub struct RemoteLogSinker {
    request_sender: BidiStreamSender<JniSinkWriterStreamRequest>,
    response_stream: BidiStreamReceiver<SinkWriterStreamResponse>,
    stream_chunk_converter: StreamChunkConverter,
    sink_writer_metrics: SinkWriterMetrics,
}

impl RemoteLogSinker {
    async fn new(

View on GitHub (pinned to 6469eb736d)