{"record":{"id":"3c43b12d78db9ccb","repo":"risingwavelabs/risingwave","slug":"bigquery-insert-error-end-of-resp-stream","errorCode":null,"errorMessage":"bigquery insert error: end of resp stream","messagePattern":"bigquery insert error: end of resp stream","errorType":"exception","errorClass":"SinkError::BigQuery","httpStatus":null,"severity":"error","filePath":"src/connector/src/sink/big_query.rs","lineNumber":853,"sourceCode":"            .map_err(|e| SinkError::BigQuery(e.into()))?\n        {\n            Some(append_rows_response) => {\n                if !append_rows_response.row_errors.is_empty() {\n                    return Err(SinkError::BigQuery(anyhow::anyhow!(\n                        \"bigquery insert error {:?}\",\n                        append_rows_response.row_errors\n                    )));\n                }\n                if let Some(google_cloud_googleapis::cloud::bigquery::storage::v1::append_rows_response::Response::Error(status)) = append_rows_response.response{\n                            return Err(SinkError::BigQuery(anyhow::anyhow!(\n                                \"bigquery insert error {:?}\",\n                                status\n                            )));\n                        }\n                yield ();\n            }\n            None => {\n                return Err(SinkError::BigQuery(anyhow::anyhow!(\n                    \"bigquery insert error: end of resp stream\",\n                )));\n            }\n        }\n    }\n}\n\nstruct StorageWriterClient {\n    #[expect(dead_code)]\n    environment: Environment,\n    request_sender: mpsc::UnboundedSender<AppendRowsRequest>,\n}\nimpl StorageWriterClient {\n    pub async fn new(\n        credentials: CredentialsFile,\n    ) -> Result<(Self, impl Stream<Item = Result<()>>)> {\n        let ts_grpc = google_cloud_auth::token::DefaultTokenSourceProvider::new_with_credentials(\n            Self::bigquery_grpc_auth_config(),","sourceCodeStart":835,"sourceCodeEnd":871,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/connector/src/sink/big_query.rs#L835-L871","documentation":"The Storage Write API AppendRows call must yield exactly one response per request. This error is returned when the response stream ended without yielding any response (None), meaning BigQuery closed or never produced a response for the append request.","triggerScenarios":"resp_to_stream awaiting the stream's .message() and receiving None before any append_rows_response - the gRPC response stream terminated prematurely after sending a batch.","commonSituations":"Network interruption or idle-timeout killing the bidirectional gRPC stream mid-batch; proxies/firewalls dropping long-lived streams; oversized batches taking too long to acknowledge; BigQuery service issues.","solutions":["Retry the sink; RisingWave fault-tolerance will replay the batch after recovery.","Check network stability between the RisingWave node and BigQuery (proxies, NAT idle timeouts); enable gRPC keepalive.","Reduce batch size / flush interval so individual AppendRows requests complete faster.","Upgrade RisingWave and the google-cloud-bigquery-storage client for stream-retry fixes; check BigQuery status for outages."],"exampleFix":null,"handlingStrategy":"retry","validationCode":null,"typeGuard":null,"tryCatchPattern":"match sink_write_result {\n    Err(SinkError::BigQuery(e)) if e.to_string().contains(\"end of resp stream\") => {\n        // transient gRPC stream close: rely on RisingWave recovery/replay or back off and retry\n    }\n    other => other?,\n}","preventionTips":["Ensure stable connectivity to bigquerystorage.googleapis.com (no aggressive idle-timeout proxies/NAT).","Keep batches modest so each AppendRows round-trip completes quickly.","Enable gRPC keepalive settings on long-running streams.","Pin a current version of the google-cloud-bigquery-storage client for retry fixes."],"tags":["bigquery","grpc","stream","network"],"backgroundTag":"request-timeout","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-14T16:17:12.679Z"}