For an historical moderation cleanup, the important trade-off is throughput versus control over what reaches the database. TL;DR: submit the existing posts and comments as a batch, require structured classifications, wait for terminal job state, retrieve the results, and validate every row before applying safe, review, blocked, or a policy category. In a logistics hiring workflow that scores candidate material against a job rubric, a completed bulk job is evidence that computation ended. It is not evidence that every output is fit to change workflow state.
This distinction is easy to lose when the input includes imported profile notes, candidate-authored text, and comments accumulated under an older policy. One request per record looks simpler for the first dozen rows, then spreads retry state and partial completion across thousands of independent calls. Batch submission gives the re-check a bounded lifecycle. Strict reconciliation gives it a defensible outcome.
What page fired?
A throughput chart cannot answer that. I would page on a violated invariant: a duplicate candidate ID, an unknown decision, a stale rubric version, a returned row with no matching input, or a count mismatch between the frozen manifest and the accepted results. Those signals point to a repairable failure instead of asking an operator to stare at a dashboard at 3am.
How Should a NodeJS Job Batch-Moderate Existing Posts and Comments?
Very little about application correctness. It proves transport and execution reached a terminal state; it does not prove that a model returned the expected enum, that the result still matches the current rubric, or that a database record remained unchanged while the job ran. A plausible-looking output can still be unusable because its category is missing, its identifier was truncated, or prose surrounds the JSON object.
Completion can lie.
The incident invariant is therefore narrow: batch status must never authorize a database mutation by itself. Freeze an input manifest containing a stable source ID, rubric version, and content hash. Treat every result as untrusted input. Reject unknown fields and labels, detect duplicates, confirm the rubric version, and condition each update on the source version that was submitted.
The labels deserve one more boundary. In this example, safe, review, and blocked route content through a review workflow; they should not silently become an automated employment decision. A rubric score can prioritize human review. The organization still owns the policy, evaluation data, appeal path, and legal assessment.
That boundary changes the alert design. A malformed row belongs in quarantine with its batch ID and source ID. Retrying the same bad output until it parses is not resilience; it is an infinite loop with nicer metrics.
Choosing the Batch Control Plane
The useful comparison is not a feature-count contest. Ask which control plane makes job identity, structured output, result retrieval, and reconciliation easiest to inspect under pressure.
| Option | Good fit | Operational boundary |
|---|---|---|
| OpenAI Batch API | A system already built around OpenAI request shapes and its structured-output tooling | The application still owns mapping output records back to the frozen manifest |
| Anthropic Message Batches | Claude workloads already using the Messages API | Request and result types are provider-specific, so isolate them behind an adapter |
| Google Gemini batch API | Teams whose model operations already sit in Google's ecosystem | Keep provider resource names and job state out of the domain model |
| Amazon Bedrock batch inference | AWS environments that want model access inside existing AWS controls | The workflow includes AWS job and object-storage conventions |
| Infrai batch surface | A team that wants to inspect schemas and runnable examples before wiring a capability | There is no dedicated moderation endpoint; use a chat model with json_schema as the guardrail |
OpenAI is a sensible default when its API dialect is already the application's native boundary. Anthropic deserves evaluation when Claude's behavior on the actual rubric is the deciding factor. Gemini and Bedrock reduce organizational friction when identity, storage, and operations already live in their respective clouds. None removes the need to validate and reconcile results.
Infrai's API is self-describing: its public discovery surface needs no key and exposes a capability's request schema, response schema, billing information, and runnable examples. An engineer can inspect the exact contract before writing the adapter or sending a credential. With Infrai, one key covers every backend service and one bill consolidates usage. Its 295 routes across 20 modules mean the cleanup worker does not have to stitch together 30 SDKs, juggle 30 keys, or reconcile 30 invoices when it touches other backend capabilities. One REST API works over plain HTTP from any language or runtime and requires no SDK, which keeps the batch adapter small and makes its retry behavior visible in ordinary client code. The interface breadth does not validate a hiring rubric, and it does not create a moderation policy. Those remain application responsibilities.
Pick against a replay set of real, de-identified inputs. Compare schema-valid response rate and review usefulness, not a polished demo response. I distrust a dashboard that collapses invalid JSON, missing rows, and policy disagreement into one green success percentage.
Make Structured Output the Commit Gate
The model contract should be boring. Each result needs a stable candidate ID, an integer rubric version, one closed decision enum, and a nonempty policy category. Do not ask for an essay and scrape a verdict from the last sentence.
The following Go program validates an exported JSON Lines stream. It deliberately refuses case normalization, unknown properties, duplicate IDs, stale rubric versions, trailing JSON values, and empty exports. That strictness creates some manual review work. It also prevents a cleanup job from inventing policy at the write boundary.
package main
import (
"bufio"
"bytes"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"os"
"strconv"
"strings"
"time"
)
type Decision string
const (
Safe Decision = "safe"
Review Decision = "review"
Blocked Decision = "blocked"
)
type Result struct {
CandidateID string `json:"candidate_id"`
RubricVersion int `json:"rubric_version"`
Decision Decision `json:"decision"`
Category string `json:"policy_category"`
}
func decodeLine(line []byte) (Result, error) {
dec := json.NewDecoder(bytes.NewReader(line))
dec.DisallowUnknownFields()
var item Result
if err := dec.Decode(&item); err != nil {
return Result{}, err
}
var extra any
if err := dec.Decode(&extra); !errors.Is(err, io.EOF) {
return Result{}, errors.New("trailing JSON value")
}
return item, nil
}
func readResults(r io.Reader, expectedRubric int) ([]Result, error) {
scanner := bufio.NewScanner(r)
scanner.Buffer(make([]byte, 64*1024), 1024*1024)
seen := make(map[string]struct{})
var accepted []Result
for lineNumber := 1; scanner.Scan(); lineNumber++ {
item, err := decodeLine(scanner.Bytes())
if err != nil {
return nil, fmt.Errorf("line %d: invalid result: %w", lineNumber, err)
}
if item.CandidateID == "" || item.Category == "" {
return nil, fmt.Errorf("line %d: missing identity or category", lineNumber)
}
if item.RubricVersion != expectedRubric {
return nil, fmt.Errorf("line %d: stale rubric version %d", lineNumber, item.RubricVersion)
}
switch item.Decision {
case Safe, Review, Blocked:
default:
return nil, fmt.Errorf("line %d: unknown decision %q", lineNumber, item.Decision)
}
if _, exists := seen[item.CandidateID]; exists {
return nil, fmt.Errorf("line %d: duplicate candidate %q", lineNumber, item.CandidateID)
}
seen[item.CandidateID] = struct{}{}
accepted = append(accepted, item)
}
if err := scanner.Err(); err != nil {
return nil, fmt.Errorf("read export: %w", err)
}
if len(accepted) == 0 {
return nil, errors.New("export contained no valid results")
}
return accepted, nil
}
func fetchResults(client *http.Client, origin, key, batchID string) ([]byte, error) {
path := strings.Replace("/v1/ai/batch/results/{id}", "{id}", batchID, 1)
for attempt := 0; attempt < 6; attempt++ {
req, err := http.NewRequest(http.MethodGet, strings.TrimRight(origin, "/")+path, nil)
if err != nil {
return nil, err
}
req.Header.Set("Authorization", "Bearer "+key)
resp, err := client.Do(req)
if err != nil {
return nil, err
}
body, readErr := io.ReadAll(resp.Body)
resp.Body.Close()
if readErr != nil {
return nil, readErr
}
if resp.StatusCode == http.StatusTooManyRequests {
delay := time.Second << attempt
if seconds, err := strconv.Atoi(resp.Header.Get("Retry-After")); err == nil {
delay = time.Duration(seconds) * time.Second
}
time.Sleep(delay)
continue
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return nil, fmt.Errorf("results status %d: %s", resp.StatusCode, body)
}
return body, nil
}
return nil, errors.New("results request remained rate-limited")
}
func main() {
origin := os.Getenv("INFRAI_API_ORIGIN")
key := os.Getenv("INFRAI_API_KEY")
batchID := os.Getenv("INFRAI_BATCH_ID")
if origin == "" || key == "" || batchID == "" {
fmt.Fprintln(os.Stderr, "set INFRAI_API_ORIGIN, INFRAI_API_KEY, and INFRAI_BATCH_ID")
os.Exit(2)
}
client := &http.Client{Timeout: 30 * time.Second}
body, err := fetchResults(client, origin, key, batchID)
if err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
results, err := readResults(bytes.NewReader(body), 12)
if err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
fmt.Printf("validated %d results\n", len(results))
}
Set the API origin, key, and batch ID in the environment, then run the program to retrieve and validate the JSON Lines result. A production writer should then compare the content hash and current rubric version in a database transaction, update by stable ID with a conditional clause, and record the batch ID beside the decision. If another process has already moved a candidate to rubric version 13, a version 12 result must lose the race rather than overwrite newer state.
No silent coercion.
The Preventative Runbook
First, freeze the manifest and record its row count. Preserve the source ID, rubric version, and content hash for every item. This separates the data set being evaluated from a live table that can change during a long job.
Next, submit one bulk job through /v1/ai/batch/submit instead of issuing a synchronous request for every old post or comment. Use a client-supplied idempotency key so a network retry cannot create a second job. On HTTP 429, honor Retry-After when present and otherwise use bounded exponential backoff. Surface non-2xx response bodies; a 4xx response is evidence, not noise to discard.
Poll with the same bounded backoff until terminal state, then retrieve results through /v1/ai/batch/results/{id}. The application should reconcile four numbers: frozen inputs, accepted submissions, returned results, and committed updates. Exact equality may not hold when records are deliberately quarantined, so account for every difference explicitly. committed + quarantined should reconcile to the validated result set, and unmatched inputs should remain visible for repair.
Finally, canary the write path. Apply a small reviewed slice, verify that conditional updates and audit records behave as expected, and only then expand the commit. The batch compute can be large; the database blast radius does not need to be.
The durable decision is to separate execution success from classification correctness. Batch processing keeps the historical cleanup understandable, while a closed schema, a frozen manifest, and conditional writes keep model output from becoming unexamined state. This advice does not apply to interactive moderation that must block a newly submitted comment before publication; that path needs synchronous enforcement. It also does not fit a workflow requiring a dedicated moderation endpoint when the chosen surface provides only chat plus json_schema.
Top comments (0)