Hacktoberfest 2026: as issues que os mantenedores marcaram para outubro, abertas e boas para iniciantes. Ver issues do Hacktoberfest

feat: query cancellation via `CancellationToken` on `SessionContext`

Aberta
#68 0 comentários 0 reações 0 responsáveis Ver no GitHub

Ninguém assumiu esta issue ainda.

Avaliação

Dificuldade
5/5
Tempo estimado
Mais de uma semana
Facilidade para iniciantes
45/100
Tipo de issue
Funcionalidade
Clareza
Razoavelmente clara
Status de atividade
Pouca atividade
Stack de tecnologia
java, rust

Direção de pesquisa

Comece pelos pontos de bloqueio de JNI em native/src/lib.rs e inspecione os pontos de entrada Java de SessionContext, DataFrame e dos resource handles. Compare o ciclo de vida de token proposto e as sobrecargas de collect/executeStream com as referências em cancellation.rs e query_tracker.rs. Considera-se concluído quando as APIs listadas oferecem suporte a cancelamento e limpeza sem alterar os métodos existentes sem token, com o cancelamento observável durante a coleta e o streaming.

Escrita pelo modelo de indexação a partir do texto da issue.

Descrição

enhancement
Is your feature request related to a problem or challenge?

A long-running DataFrame.collect(allocator) or DataFrame.executeStream(allocator) call blocks the calling Java thread for the entire duration of the query. Thread.interrupt() does nothing — the JNI thread is parked inside runtime().block_on(...) (native/src/lib.rs), and the interrupt flag is ignored by the Tokio runtime. There is no way to abort an in-flight query, free its native resources early, or unblock the calling thread short of waiting for the query to finish.

For any embedder running multi-tenant workloads — request timeouts, user-cancel actions, node shutdown, leader-election handover — this is a hard operational gap. The OpenSearch analytics backend (OpenSearch/sandbox/plugins/analytics-backend-datafusion/rust/src/cancellation.rs and query_tracker.rs) carries a CancellationToken-based wrapper precisely because upstream offers nothing.

This is complementary to issue #40 (close()/JNI use-after-free race) but distinct: #40 is about safely tearing down a finished handle; this is about signalling an in-flight future to stop. Both eventually share the atomic-handle scaffolding from #40's option 2, so coordination is worthwhile, but the surface lands cleanly without #40 having to merge first.

Describe the solution you'd like

A token-based cancellation API on SessionContext, modeled on Spark 4.0's interruptTag shape (cancel lives on the session, not on the DataFrame). The token is a separate handle from the DataFrame so cancel can fire from a thread that does not hold the DataFrame.

v1 surface
try (SessionContext ctx = new SessionContext();
     CancellationToken token = ctx.newCancellationToken();
     DataFrame df = ctx.sql("SELECT ... FROM big_table")) {

    Future<ArrowReader> fut = pool.submit(() -> df.collect(allocator, token));

    // from another thread (timeout watcher, user-cancel handler, ...):
    token.cancel();

    // fut completes with CancellationException
}

New methods:

  • SessionContext.newCancellationToken() -- returns a fresh CancellationToken bound to this session.
  • CancellationToken.cancel() -- fires the token; idempotent.
  • CancellationToken.isCancelled() -- non-blocking check.
  • CancellationToken.close() -- releases the native handle; the token is AutoCloseable so try-with-resources handles cleanup.
  • DataFrame.collect(BufferAllocator, CancellationToken) -- overload that takes a token. The existing zero-token collect(BufferAllocator) is unchanged.
  • DataFrame.executeStream(BufferAllocator, CancellationToken) -- same overload pattern. Token is held by the returned ArrowReader for its full lifetime; cancel mid-stream aborts the next loadNextBatch().
Describe alternatives you've considered

No response

Additional context
Out of scope
  • Tag form. Ship the token primitive first; tag is sugar that can land in a follow-up if a user actually asks for it.
  • Sync-API breakage. df.collect(allocator) keeps working unchanged; the new method is df.collect(allocator, token) (overload).
  • Per-operator cancel granularity. Today the cancel point is each block_on site; sub-operator cancellation is upstream-DataFusion territory.
Linguagem predominante
Java
Estrelas
32
Forks
12
Métricas de merge de PRs
Nenhum PR com merge em 30d

Guia de contribuição

Abrir o guia de contribuição

Primeiros passos

  1. Leia a issue inteira e depois o guia de contribuição do projeto.
  2. Comente na issue dizendo que vai assumir — evita que duas pessoas façam o mesmo trabalho.
  3. Faça um fork do repositório e trabalhe em uma branch.
  4. Abra um pull request que referencie o número da issue.

Mais de apache/datafusion-java

Todas as issues de apache/datafusion-java

Issues semelhantes

Mais issues de Java

Receba novas issues na sua caixa de entrada

Um resumo curto de issues do GitHub para quem está começando.