Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -186,6 +186,17 @@ publishing {
}
}

// Central rejects any component missing sources/javadoc, and one bad component fails the whole
// deployment bundle — which would take chalk-java down with it. Reuse the unshaded jars: relocation
// is a bytecode rewrite, so the sources are identical and the public ai.chalk.* API is unrelocated.
// Deferred to afterEvaluate because the vanniktech plugin registers both tasks after this file runs.
afterEvaluate {
publishing.publications.named('shaded') {
artifact(tasks.named('sourcesJar')) { classifier = 'sources' }
artifact(tasks.named('mavenPlainJavadocJar')) { classifier = 'javadoc' }
}
}

// With two publications (maven + shaded), both sign tasks write their .asc into build/libs, and each
// publish task reads from build/libs — so Gradle's strict validation flags an implicit dependency
// between publish and the *other* publication's sign task and fails the build. Order all publish
Expand Down
8 changes: 7 additions & 1 deletion docs/flink-sink.md
Original file line number Diff line number Diff line change
Expand Up @@ -105,12 +105,18 @@ Each element maps to `Map<feature_fqn, value>`. Rules:
| `batchSize` | `1000` | Flush when this many rows are buffered. |
| `flushInterval` | `5s` | Flush a non-empty buffer at least this often. |
| `uploadTimeout` | `30s` | Per-call deadline. |
| `maxRetries` | `3` | Retries per batch on transient failure (exponential backoff). |
| `retryTimeout` | `120s` | Total wall-clock budget for retrying one batch's transient failures. The primary retry bound — sized to outlast a routine query-server rollout. |
| `maxRetries` | unlimited | Optional hard cap on retries per batch. Off by default so `retryTimeout` governs; set it to stop after a fixed number of attempts regardless of remaining budget. |
| `retryBackoff` | `500ms` | Base backoff, doubled per attempt (capped 30s). |
| `failOnUploadErrors` | `true` | `true`: fail the sink on engine data-level errors (replay). `false`: log a WARN and **drop** the rejected rows, continuing. |

> **Poison-pill note:** engine *data-level* errors (e.g. a bad feature type) are deterministic — with `failOnUploadErrors=true`, replaying the same batch after a checkpoint restart hits the same error, wedging the pipeline in a restart loop. Fix such rows upstream, or set `failOnUploadErrors=false` to drop them. (Transient transport errors are handled separately by retry + Flink restart and do make progress.)

> **Retry budget:** retries block the task thread, so `retryTimeout` is also the worst-case stall for
> that subtask — keep it below the checkpoint timeout. An outage longer than the budget still surfaces
> as a task failure and a restart from the last checkpoint; that is at-least-once working as intended,
> not data loss.

## Write targets

`ChalkSinkConfig` exposes `writeOnline` (default `true`), `writeOffline`, and `updateMataggs`
Expand Down
61 changes: 43 additions & 18 deletions src/main/java/ai/chalk/flink/core/ChalkFeatureUploader.java
Original file line number Diff line number Diff line change
Expand Up @@ -10,14 +10,20 @@
import java.util.LinkedHashSet;
import java.util.List;
import java.util.Map;
import java.util.concurrent.TimeUnit;

/**
* Flink-free batching engine for Chalk {@code upload_features}. Accumulates feature rows and flushes
* them as a single columnar upload when the buffer fills (see {@link #add}), or when the owning sink
* asks (checkpoint / close). Retries <em>transient</em> failures (gRPC UNAVAILABLE / DEADLINE_EXCEEDED
* / RESOURCE_EXHAUSTED / ABORTED) with exponential backoff; non-transient failures (auth, invalid
* argument) are thrown immediately. On exhaustion it throws, letting the caller (a Flink sink) fail
* and replay from the last checkpoint — at-least-once, made safe by upsert semantics.
* / RESOURCE_EXHAUSTED / ABORTED) with exponential backoff until the {@code retryTimeout} budget is
* spent; non-transient failures (auth, invalid argument) are thrown immediately. On exhaustion it
* throws, letting the caller (a Flink sink) fail and replay from the last checkpoint — at-least-once,
* made safe by upsert semantics.
*
* <p>The budget is <em>time</em>, not attempts, so it survives a backend rollout of a known duration
* regardless of how the backoff is tuned. Retries run inline on the task thread, so the budget is
* also the worst-case subtask stall.
*
* <p><b>Flush latency</b> is bounded by {@code min(batchSize reached, checkpoint interval)}. The
* configured flush interval is a best-effort upper bound evaluated when the next element arrives; on
Expand Down Expand Up @@ -98,26 +104,39 @@ private boolean intervalElapsed() {
}

private UploadOutcome uploadWithRetry(Map<String, List<?>> columnar) {
int attempts = config.maxRetries() + 1;
// nanoTime, not currentTimeMillis: the budget must not shift under an NTP step.
long start = System.nanoTime();
long budgetNanos = TimeUnit.MILLISECONDS.toNanos(config.retryTimeoutMillis());
Exception last = null;
for (int attempt = 0; attempt < attempts; attempt++) {
int attempts = 0;
while (true) {
try {
return client.upload(columnar);
} catch (Exception e) {
last = e;
attempts++;
// Fail fast on non-transient errors (auth, invalid argument, ...): retrying them
// just burns backoff and then triggers a Flink restart storm on the same error.
if (attempt >= attempts - 1 || !isRetryable(e)) {
// just burns the budget and then triggers a Flink restart storm on the same error.
if (!isRetryable(e) || attempts > config.maxRetries()) {
break;
}
long backoff = backoffMillis(attempts - 1);
// Stop once the next sleep would outlast the budget: sleeping past the deadline
// delays the restart without buying another attempt.
long remaining = TimeUnit.NANOSECONDS.toMillis(budgetNanos - (System.nanoTime() - start));
if (backoff >= remaining) {
break;
}
LOG.warn("transient upload_features failure (attempt {}/{}), retrying: {}",
attempt + 1, attempts, e.toString());
sleepBackoff(attempt);
LOG.warn("transient upload_features failure (attempt {}, {}ms of retry budget left), "
+ "retrying in {}ms: {}", attempts, remaining, backoff, e.toString());
sleep(backoff);
}
}
boolean retryable = isRetryable(last);
long elapsedMillis = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - start);
throw new ChalkUploadException(
"upload_features failed" + (retryable ? " after " + attempts + " attempt(s)"
"upload_features failed" + (isRetryable(last)
? " after " + attempts + " attempt(s) over " + elapsedMillis + "ms"
+ " (retryTimeout=" + config.retryTimeoutMillis() + "ms)"
: " with a non-retryable error"), last);
}

Expand All @@ -135,13 +154,19 @@ private static boolean isRetryable(Throwable t) {
return false;
}

private void sleepBackoff(int attempt) {
// Exponential backoff, capped at 30s to bound restart latency. Cap the shift and detect
// overflow (a huge maxRetries could otherwise wrap the product negative -> bad Thread.sleep).
long scaled = config.retryBackoffMillis() << Math.min(attempt, 20);
long backoff = (scaled < 0) ? 30_000L : Math.min(scaled, 30_000L);
/**
* Exponential backoff for the given zero-based retry index, capped at 30s so a long budget is
* spent on many attempts rather than a few enormous sleeps. The shift is capped and overflow
* detected: an unbounded attempt count could otherwise wrap the product negative.
*/
private long backoffMillis(int retryIndex) {
long scaled = config.retryBackoffMillis() << Math.min(retryIndex, 20);
return (scaled < 0) ? 30_000L : Math.min(scaled, 30_000L);
}

private void sleep(long millis) {
try {
Thread.sleep(backoff);
Thread.sleep(millis);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new ChalkUploadException("interrupted while backing off before upload retry", ie);
Expand Down
31 changes: 28 additions & 3 deletions src/main/java/ai/chalk/flink/core/ChalkSinkConfig.java
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ public final class ChalkSinkConfig implements Serializable {
private final long uploadTimeoutMillis;
private final int maxRetries;
private final long retryBackoffMillis;
private final long retryTimeoutMillis;
private final boolean failOnUploadErrors;

private final boolean writeOnline;
Expand All @@ -46,6 +47,7 @@ private ChalkSinkConfig(Builder b) {
this.uploadTimeoutMillis = b.uploadTimeoutMillis;
this.maxRetries = b.maxRetries;
this.retryBackoffMillis = b.retryBackoffMillis;
this.retryTimeoutMillis = b.retryTimeoutMillis;
this.failOnUploadErrors = b.failOnUploadErrors;
this.writeOnline = b.writeOnline;
this.writeOffline = b.writeOffline;
Expand All @@ -63,6 +65,8 @@ private ChalkSinkConfig(Builder b) {
public Duration uploadTimeout() { return Duration.ofMillis(uploadTimeoutMillis); }
public int maxRetries() { return maxRetries; }
public long retryBackoffMillis() { return retryBackoffMillis; }
public long retryTimeoutMillis() { return retryTimeoutMillis; }
public Duration retryTimeout() { return Duration.ofMillis(retryTimeoutMillis); }
public boolean failOnUploadErrors() { return failOnUploadErrors; }
public boolean writeOnline() { return writeOnline; }
public boolean writeOffline() { return writeOffline; }
Expand All @@ -82,8 +86,12 @@ public static final class Builder {
private int batchSize = 1_000;
private long flushIntervalMillis = 5_000;
private long uploadTimeoutMillis = 30_000;
private int maxRetries = 3;
// Unlimited by default: retryTimeout is the real bound. A fixed attempt count can't express
// "survive a backend rollout" — with the default backoff, 3 retries covered only ~3.5s, which
// is shorter than a routine query-server deploy and caused Flink restarts in production.
private int maxRetries = Integer.MAX_VALUE;
private long retryBackoffMillis = 500;
private long retryTimeoutMillis = 120_000;
private boolean failOnUploadErrors = true;
private boolean writeOnline = true;
private boolean writeOffline = false;
Expand Down Expand Up @@ -116,12 +124,26 @@ public static final class Builder {
/** Per-call upload deadline. Default 30s. */
public Builder uploadTimeout(Duration v) { this.uploadTimeoutMillis = v.toMillis(); return this; }

/** Retries per batch on transient failures (in addition to the first attempt). Default 3. */
/**
* Optional hard cap on retries per batch, in addition to the first attempt. Unlimited by
* default — {@link #retryTimeout} is the primary bound. Set this only to stop retrying after
* a fixed number of attempts regardless of how much time budget remains.
*/
public Builder maxRetries(int v) { this.maxRetries = v; return this; }

/** Base backoff between retries, in millis (doubled each attempt). Default 500. */
/** Base backoff between retries, in millis (doubled each attempt, capped at 30s). Default 500. */
public Builder retryBackoff(Duration v) { this.retryBackoffMillis = v.toMillis(); return this; }

/**
* Total wall-clock budget for retrying one batch's transient failures. Default 120s, chosen to
* outlast a routine query-server rollout: a shorter budget surfaces the outage as a Flink task
* failure and a restart from the last checkpoint.
*
* <p>Retries block the task thread, so this is also the worst-case stall for the subtask.
* Keep it below the checkpoint timeout.
*/
public Builder retryTimeout(Duration v) { this.retryTimeoutMillis = v.toMillis(); return this; }

/**
* When true (default), a batch whose response carries engine data-level errors fails the sink.
*
Expand Down Expand Up @@ -171,6 +193,9 @@ public ChalkSinkConfig build() {
if (retryBackoffMillis < 0) {
throw new IllegalArgumentException("retryBackoff must be >= 0");
}
if (retryTimeoutMillis < 0) {
throw new IllegalArgumentException("retryTimeout must be >= 0");
}
// Fail fast at job-construction time (not later on the TaskManager) for write targets the
// pinned chalk-java can't honor. See README "Write targets".
if (writeOffline || updateMataggs || !writeOnline) {
Expand Down
48 changes: 48 additions & 0 deletions src/test/java/ai/chalk/flink/core/ChalkFeatureUploaderTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,54 @@ void throwsAfterRetriesExhausted() {
assertEquals(3, fake.attempts.get()); // 1 initial + 2 retries
}

@Test
void stopsRetryingWhenTimeBudgetSpent() {
FakeClient fake = new FakeClient();
fake.failFirstN = Integer.MAX_VALUE; // never recovers
ChalkFeatureUploader uploader = new ChalkFeatureUploader(
config(ChalkSinkConfig.builder().batchSize(1)
.retryTimeout(Duration.ofMillis(300))
.retryBackoff(Duration.ofMillis(50))), fake);

// maxRetries is unlimited by default, so the time budget is the only thing that can stop this.
long start = System.nanoTime();
ChalkUploadException e = assertThrows(ChalkUploadException.class,
() -> uploader.add(Map.of("user.id", 1L)));
long elapsedMillis = (System.nanoTime() - start) / 1_000_000;

assertTrue(e.getMessage().contains("retryTimeout"), e.getMessage());
// Never sleeps past the deadline; generous ceiling so a slow CI box can't flake this.
assertTrue(elapsedMillis < 3_000, "took " + elapsedMillis + "ms, expected to stop near 300ms");
}

@Test
void timeBudgetAllowsMoreAttemptsThanTheOldFixedCap() {
FakeClient fake = new FakeClient();
fake.failFirstN = Integer.MAX_VALUE;
ChalkFeatureUploader uploader = new ChalkFeatureUploader(
config(ChalkSinkConfig.builder().batchSize(1)
.retryTimeout(Duration.ofMillis(1_000))
.retryBackoff(Duration.ofMillis(1))), fake);

assertThrows(ChalkUploadException.class, () -> uploader.add(Map.of("user.id", 1L)));
// The old behaviour was a hard 4 attempts (maxRetries=3 + 1) regardless of remaining budget;
// with a 1s budget and 1ms base backoff there is room for many more.
assertTrue(fake.attempts.get() > 4, "only " + fake.attempts.get() + " attempts");
}

@Test
void explicitMaxRetriesStillCapsBeforeBudget() {
FakeClient fake = new FakeClient();
fake.failFirstN = Integer.MAX_VALUE;
ChalkFeatureUploader uploader = new ChalkFeatureUploader(
config(ChalkSinkConfig.builder().batchSize(1).maxRetries(2)
.retryTimeout(Duration.ofMinutes(5)) // budget far larger than the cap
.retryBackoff(Duration.ofMillis(1))), fake);

assertThrows(ChalkUploadException.class, () -> uploader.add(Map.of("user.id", 1L)));
assertEquals(3, fake.attempts.get()); // 1 initial + 2 retries, budget untouched
}

@Test
void failsFastOnNonRetryableError() {
FakeClient fake = new FakeClient();
Expand Down
10 changes: 10 additions & 0 deletions src/test/java/ai/chalk/flink/core/ChalkSinkConfigTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import java.time.Duration;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;

class ChalkSinkConfigTest {
Expand All @@ -28,11 +29,20 @@ void rejectsUnsupportedWriteTargetsAtConstruction() {
@Test
void rejectsInvalidDurationsAndCounts() {
assertThrows(IllegalArgumentException.class, () -> base().retryBackoff(Duration.ofMillis(-1)).build());
assertThrows(IllegalArgumentException.class, () -> base().retryTimeout(Duration.ofMillis(-1)).build());
assertThrows(IllegalArgumentException.class, () -> base().flushInterval(Duration.ofMillis(-1)).build());
assertThrows(IllegalArgumentException.class, () -> base().uploadTimeout(Duration.ZERO).build());
assertThrows(IllegalArgumentException.class, () -> base().batchSize(0).build());
}

@Test
void retryBudgetDefaultsToTimeNotAttempts() {
ChalkSinkConfig c = base().build();
assertEquals(Duration.ofSeconds(120), c.retryTimeout());
// Unlimited by default so retryTimeout is the effective bound.
assertEquals(Integer.MAX_VALUE, c.maxRetries());
}

@Test
void rejectsMissingCredentials() {
assertThrows(IllegalArgumentException.class,
Expand Down
Loading