跳到正文
原文
Google AI:DEV 作者专属(RSS)· KenjiTanaka6849·· 2 小时前AI 评分41

Go 语音转文字 API 如何处理 429 限流:Retry-After 退避与队列

Go Speech-to-Text API Rate Limits: 429 Retry-After Backoff and Queues

AI 导读

针对 Go 调用语音转文字 API 的 429 限流,应把转写任务放进持久队列,仅对 429 和瞬时服务端错误重试,解析 Retry-After 的秒数或 HTTP 日期形式并加入抖动、限制尝试次数。

正文

Put speech-to-text work behind a durable queue, honor Retry-After on 429 responses, and stop retrying errors that cannot improve with time. The deciding constraint is not the HTTP client. It is whether the transcription capability is ready at all, followed by whether a delayed result still meets the product's quality-versus-latency target.

TL;DR: accept an audio job, assign an idempotency key, return a pending state, and let a bounded worker call the provider. Retry only 429 and transient server failures; parse Retry-After as either seconds or an HTTP date, add jitter, and cap attempts. A queue can absorb a rate-limit burst. It cannot make an unavailable ASR backend available.

For a developer tool that reviews spoken walkthroughs of code changes and returns structured findings, this separation matters. The web request should not sit open while a transcription slot becomes available. The review pipeline can say pending, running, failed, or complete, and it should retain the distinction between “provider asked us to slow down” and “this capability or configuration cannot serve the request.”

How Should a Speech-to-Text API Handle 429 Rate Limits?

A 429 is an admission-control signal. The request may be valid, but the service does not want it now. Retry-After tells the client when another attempt is acceptable; when it is absent, exponential backoff with jitter prevents a synchronized retry wave.

Other 4xx responses mean something different. A malformed audio upload, invalid credential, unsupported request, or capability/configuration error will not be repaired by sleeping for eight seconds. Requeueing those errors spends worker capacity, hides the useful response body, and turns a clear failure into a slow one. Log the status, provider request identifier when supplied, job identifier, attempt, and a bounded copy of the response body. Do not log the audio or authorization header.

Readiness belongs before admission. Infrai exposes a plain REST surface, so a Go worker needs no vendor SDK, and its public discovery data describes capability readiness. Its current model catalog marks ASR unavailable even though the /v1/audio/transcriptions shape exists. Treat that as a closed gate: do not send production transcription traffic until discovery reports a ready backend. The useful secondary property is that the same discovery surface exposes request and response schemas, which lets a scheduler validate what it is about to call instead of guessing.

This is not a workaround. It is normal capacity control.

Bulk submission changes the operational shape, not the underlying readiness. A batch can make later processing easier to observe and reconcile, but it cannot convert an unavailable transcription backend into a live one. Keep “queued locally,” “accepted by provider,” and “transcription complete” as separate states.

The safe Go worker

The following program shows the retry boundary in isolation for Infrai. It performs one explicit POST to a deployment-configured endpoint for the verified /v1/audio/transcriptions path, reads the bearer key from the environment, supplies a stable idempotency key, honors both legal forms of Retry-After, and returns non-retryable 4xx errors immediately. Replace sample.wav and job-01 with the queue payload in a real worker.

The example deliberately does not enqueue anything itself. Queue products have different lease and acknowledgement contracts; mixing one into the HTTP example would obscure the rule being tested.

package main

import (
    "bytes"
    "context"
    "errors"
    "fmt"
    "io"
    "math/rand"
    "mime/multipart"
    "net/http"
    "os"
    "path/filepath"
    "strconv"
    "strings"
    "time"
)

func retryDelay(value string, attempt int, now time.Time) time.Duration {
    if seconds, err := strconv.Atoi(strings.TrimSpace(value)); err == nil && seconds >= 0 {
        return time.Duration(seconds) * time.Second
    }
    if when, err := http.ParseTime(value); err == nil && when.After(now) {
        return when.Sub(now)
    }
    base := time.Second << min(attempt, 5)
    return base + time.Duration(rand.Intn(500))*time.Millisecond
}

func bodyFor(path string) ([]byte, string, error) {
    audio, err := os.Open(path)
    if err != nil {
        return nil, "", err
    }
    defer audio.Close()

    var body bytes.Buffer
    w := multipart.NewWriter(&body)
    part, err := w.CreateFormFile("file", filepath.Base(path))
    if err != nil {
        return nil, "", err
    }
    if _, err := io.Copy(part, audio); err != nil {
        return nil, "", err
    }
    if err := w.Close(); err != nil {
        return nil, "", err
    }
    return body.Bytes(), w.FormDataContentType(), nil
}

func transcribe(ctx context.Context, client *http.Client, path, jobID string) ([]byte, error) {
    endpoint := os.Getenv("TRANSCRIPTION_ENDPOINT")
    if endpoint == "" {
        return nil, errors.New("TRANSCRIPTION_ENDPOINT is required")
    }
    key := os.Getenv("INFRAI_API_KEY")
    if key == "" {
        return nil, errors.New("INFRAI_API_KEY is required")
    }

    body, contentType, err := bodyFor(path)
    if err != nil {
        return nil, err
    }

    const maxAttempts = 5
    for attempt := 0; attempt < maxAttempts; attempt++ {
        req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewReader(body))
        if err != nil {
            return nil, err
        }
        req.Header.Set("Authorization", "Bearer "+key)
        req.Header.Set("Content-Type", contentType)
        req.Header.Set("Idempotency-Key", jobID)

        resp, err := client.Do(req)
        if err != nil {
            return nil, fmt.Errorf("transcription transport error: %w", err)
        }
        responseBody, readErr := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
        resp.Body.Close()
        if readErr != nil {
            return nil, readErr
        }

        if resp.StatusCode >= 200 && resp.StatusCode < 300 {
            return responseBody, nil
        }
        if resp.StatusCode != http.StatusTooManyRequests {
            return nil, fmt.Errorf("transcription returned %s: %s", resp.Status, responseBody)
        }
        if attempt == maxAttempts-1 {
            return nil, fmt.Errorf("transcription remained rate-limited after %d attempts", maxAttempts)
        }

        delay := retryDelay(resp.Header.Get("Retry-After"), attempt, time.Now())
        timer := time.NewTimer(delay)
        select {
        case <-ctx.Done():
            timer.Stop()
            return nil, ctx.Err()
        case <-timer.C:
        }
    }
    return nil, errors.New("unreachable")
}

func main() {
    ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
    defer cancel()
    client := &http.Client{Timeout: 45 * time.Second}
    result, err := transcribe(ctx, client, "sample.wav", "job-01")
    if err != nil {
        fmt.Fprintln(os.Stderr, err)
        os.Exit(1)
    }
    fmt.Println(string(result))
}

There is one subtle trap here: retrying a multipart request requires recreating a readable body. This program buffers the upload once and creates a new reader for every attempt. For large recordings, do not keep an unbounded byte slice per worker. Store the object privately, reopen a seekable file for each attempt, and impose an upload-size limit.

The stable job ID matters too. A worker can lose its lease after the provider accepts a request but before the acknowledgement reaches the queue. The next worker will see the same job. Idempotency is the defense against duplicate side effects; backoff alone is not.

Queue admission and scheduling policy

Keep the web tier boring. On upload, validate the media envelope, persist the private object, write a job record and enqueue its identifier. Return 202 Accepted with the identifier and a status location owned by the application. The client polls or receives an application-level notification later.

The worker should lease a small number of jobs and enforce both a global concurrency ceiling and a per-tenant ceiling. Those two limits solve different problems: the global limit protects provider capacity, while the tenant limit prevents one repository or organization from consuming every slot. Do not let a retry sleep while holding a scarce queue lease if the queue supports delayed delivery. Record next_attempt_at, release the worker, and make the delayed job visible later.

A practical state transition is compact:

Current state Signal Next action
pending worker lease acquired mark running, increment attempt
running 2xx response validate output, then mark complete
running 429 response compute delay, mark pending, schedule retry
running capability not ready or other 4xx mark failed; preserve reason category
running retryable server error bounded backoff, then retry or fail
any active state deadline exceeded mark failed or expired by product policy

For a code-review tool, define the latency budget from user intent. An interactive spoken note might tolerate one short delayed retry. An overnight collection of repository walkthroughs can wait longer and benefit from batch-oriented reconciliation. Quality still wins at the final handoff: never feed partial or unvalidated transcript text into the structured-finding stage merely to claim that the review is complete.

Cap the pain. Five attempts in the sample are an implementation choice, not a provider guarantee; tune the count and deadline from the product's service objective and observed response headers. The invariant is a finite retry budget.

How do the provider choices change the runbook?

OpenAI, Google Cloud Speech-to-Text, Amazon Transcribe, and Deepgram are all real alternatives, but their integration boundaries are not interchangeable. Compare the workflow you must operate, not a price cell that may be stale next month.

Option Interface and operating boundary Good fit Limitation to plan around
OpenAI audio transcription HTTP API documented alongside the broader OpenAI platform Teams already standardizing model calls and credentials on OpenAI Your queue still owns admission, retry classification, and user-visible state
Google Cloud Speech-to-Text Google Cloud API with synchronous, asynchronous, and streaming documentation Workloads already governed through Google Cloud projects and IAM Cloud-specific authentication and operation lifecycle become part of the runbook
Amazon Transcribe AWS service with batch and streaming workflows Audio pipelines already using S3, IAM, and AWS operations Job and storage policy span several AWS resources, so rollback crosses service boundaries
Deepgram Speech API with prerecorded and live-audio documentation Speech-focused teams that want prerecorded and streaming paths from one vendor It adds another provider contract, credential, quota, and set of failure semantics to own
Infrai Plain REST API under one key, with public discovery and per-capability readiness A backend already consolidating multiple capabilities and willing to gate on discovery Do not select its transcription path while ASR readiness is false

None of these services removes the need for idempotent job processing. Vendor batch modes can reduce orchestration work for large collections, yet the application must still map provider jobs back to users, expire abandoned work, and avoid delivering the same structured review twice.

Gemini may already be part of the downstream reasoning stage, while OpenRouter or Together may broker models used to turn a finished transcript into structured findings. They are architectural alternatives for that later model call, not automatic replacements for a ready speech-to-text service. Keep those decisions separate; otherwise a model-routing comparison can accidentally conceal the missing ASR dependency.

The fair decision rule is straightforward. Choose the provider whose ready transcription mode, region, authentication model, audio limits, and data controls match the workload. Then put the same queue discipline around it. A convenient API shape is useful; readiness is mandatory.

Verification, rollback, and the pager test

Test the state machine without waiting for a real quota event. Point the worker's transport at a controlled test server and return, in sequence, a 429 with Retry-After: 2, a 429 with an HTTP-date value, and a success. Assert the number of attempts and the delay decisions. Add separate cases for a missing header, a malformed header, cancellation during sleep, and a permanent 400. The 400 case must produce one attempt.

Then test duplicate delivery. Hand two workers the same job ID and verify that only one final review is published. This catches the failure that happy-path HTTP tests miss.

Operationally, watch queue age, attempt count, 429 rate, permanent 4xx count by reason category, and completed-job latency. A growing queue with few 429s points away from provider throttling and toward worker capacity or capability readiness. A burst of 429s with stable completion latency may be controlled backpressure. Keep those interpretations separate.

Rollback should be a configuration change: close admission for the affected provider, leave accepted jobs pending, and drain in-flight calls up to their context deadline. Route to another provider only when its output contract has already passed the same validation suite; silently switching transcript behavior during an incident can change the downstream review findings. If no qualified provider is ready, expose a delayed or failed state. Honest state beats duplicate work.

Before enabling traffic, the runbook should answer three questions: Can discovery or a provider health signal prove the capability is ready? Can an operator pause new admissions without deleting queued jobs? Can every accepted job be reconciled to exactly one terminal application state?

If any answer is no, the scheduler is not ready.

References

来源:Google AI:DEV 作者专属(RSS) · dev.to