Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -25,11 +25,16 @@
import org.apache.dolphinscheduler.api.utils.Result;
import org.apache.dolphinscheduler.common.constants.Constants;
import org.apache.dolphinscheduler.dao.entity.ResponseTaskLog;
import org.apache.dolphinscheduler.dao.entity.TaskInstance;
import org.apache.dolphinscheduler.dao.entity.User;

import javax.servlet.http.HttpServletRequest;
import javax.servlet.http.HttpServletResponse;

import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.http.HttpHeaders;
import org.springframework.http.HttpStatus;
import org.springframework.http.MediaType;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.GetMapping;
import org.springframework.web.bind.annotation.RequestAttribute;
Expand All @@ -38,6 +43,9 @@
import org.springframework.web.bind.annotation.ResponseBody;
import org.springframework.web.bind.annotation.ResponseStatus;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.web.context.request.async.StandardServletAsyncWebRequest;
import org.springframework.web.context.request.async.WebAsyncUtils;
import org.springframework.web.servlet.mvc.method.annotation.StreamingResponseBody;

import io.swagger.v3.oas.annotations.Operation;
import io.swagger.v3.oas.annotations.Parameter;
Expand All @@ -50,6 +58,15 @@
@RequestMapping("/log")
public class LoggerController extends BaseController {

/**
* Endpoint-scoped async timeout for the log download: {@link StreamingResponseBody} runs on
* the async request path where the servlet container's default timeout (30s in Jetty/Tomcat)
* SILENTLY truncates any download that takes longer — large log downloads regularly exceed
* it. Scoped to THIS endpoint instead of the global {@code spring.mvc.async.request-timeout}
* so that no other endpoint ever inherits a 1-hour async timeout.
*/
private static final long LOG_DOWNLOAD_ASYNC_TIMEOUT_MILLIS = 60 * 60 * 1000L;

@Autowired
private LoggerService loggerService;

Expand Down Expand Up @@ -92,14 +109,30 @@ public Result<ResponseTaskLog> queryLog(@Parameter(hidden = true) @RequestAttrib
@GetMapping(value = "/download-log")
@ResponseBody
@ApiException(DOWNLOAD_TASK_INSTANCE_LOG_FILE_ERROR)
public ResponseEntity downloadTaskLog(@Parameter(hidden = true) @RequestAttribute(value = Constants.SESSION_USER) User loginUser,
@RequestParam(value = "taskInstanceId") int taskInstanceId) {
byte[] logBytes = loggerService.getLogBytes(loginUser, taskInstanceId);
public ResponseEntity<StreamingResponseBody> downloadTaskLog(
@Parameter(hidden = true) @RequestAttribute(value = Constants.SESSION_USER) User loginUser,
@RequestParam(value = "taskInstanceId") int taskInstanceId,
@Parameter(hidden = true) HttpServletRequest request,
@Parameter(hidden = true) HttpServletResponse response) {
// Sync auth check — throws ServiceException BEFORE response is committed,
// so @ApiException can still return a proper JSON error
final TaskInstance taskInstance = loggerService.checkDownloadLogAuth(loginUser, taskInstanceId);
// Scope a long async timeout to THIS request only (after the auth check, so failures
// stay on the non-async JSON error path). The servlet container's default async
// timeout (30s in Jetty/Tomcat) silently truncates StreamingResponseBody downloads
// that take longer. The WebAsyncManager is request-scoped: installing a custom
// AsyncWebRequest here affects only this endpoint — unlike the global
// spring.mvc.async.request-timeout.
final StandardServletAsyncWebRequest asyncWebRequest = new StandardServletAsyncWebRequest(request, response);
asyncWebRequest.setTimeout(LOG_DOWNLOAD_ASYNC_TIMEOUT_MILLIS);
WebAsyncUtils.getAsyncManager(request).setAsyncWebRequest(asyncWebRequest);
final StreamingResponseBody body = outputStream -> loggerService.streamLogBytes(taskInstance, outputStream);
return ResponseEntity
.ok()
.contentType(MediaType.APPLICATION_OCTET_STREAM)
.header(HttpHeaders.CONTENT_DISPOSITION,
"attachment; filename=\"" + System.currentTimeMillis() + ".log" + "\"")
.body(logBytes);
.body(body);
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -33,18 +33,6 @@
@Slf4j
public class LocalLogClient {

/**
* Download the complete log of a task instance.
* This method is used to retrieve all log information from the start to the end of a task instance,
* suitable for scenarios where a complete log record is required.
*
* @param taskInstance The task instance object, containing information needed to retrieve the log.
* @return The complete log file download response of the task instance, including log content and metadata.
*/
public TaskInstanceLogFileDownloadResponse getWholeLog(TaskInstance taskInstance) {
return getLocalWholeLog(taskInstance);
}

/**
* Query a portion of the log of a task instance.
* This method is used to query log information of a task instance in a paginated manner,
Expand All @@ -59,11 +47,14 @@ public TaskInstanceLogPageQueryResponse getPartLog(TaskInstance taskInstance, in
return getLocalPartLog(taskInstance, skipLineNum, limit);
}

private TaskInstanceLogFileDownloadResponse getLocalWholeLog(TaskInstance taskInstance) {
TaskInstanceLogFileDownloadRequest request = new TaskInstanceLogFileDownloadRequest(
taskInstance.getId(),
taskInstance.getLogPath());
return getProxyLogService(taskInstance).getTaskInstanceWholeLogFileBytes(request);
/**
* Fetch a single bounded chunk of the task instance log from the worker via chunked RPC.
*/
public TaskInstanceLogFileDownloadResponse getLogChunk(final TaskInstance taskInstance,
final long offset, final int length) {
final TaskInstanceLogFileDownloadRequest request = new TaskInstanceLogFileDownloadRequest(
taskInstance.getId(), taskInstance.getLogPath(), offset, length);
return getProxyLogService(taskInstance).getTaskInstanceLogFileChunk(request);
}

private TaskInstanceLogPageQueryResponse getLocalPartLog(TaskInstance taskInstance, int skipLineNum,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,13 +18,19 @@
package org.apache.dolphinscheduler.api.executor.logging;

import org.apache.dolphinscheduler.dao.entity.TaskInstance;
import org.apache.dolphinscheduler.extract.base.exception.MethodInvocationException;
import org.apache.dolphinscheduler.extract.common.transportor.LogResponseStatus;
import org.apache.dolphinscheduler.extract.common.transportor.TaskInstanceLogFileDownloadResponse;
import org.apache.dolphinscheduler.extract.common.transportor.TaskInstanceLogPageQueryResponse;
import org.apache.dolphinscheduler.plugin.task.api.utils.TaskTypeUtils;
import org.apache.dolphinscheduler.registry.api.RegistryClient;
import org.apache.dolphinscheduler.registry.api.enums.RegistryNodeType;

import org.apache.commons.lang3.exception.ExceptionUtils;

import java.io.IOException;
import java.io.OutputStream;

import lombok.extern.slf4j.Slf4j;

import org.springframework.beans.factory.annotation.Autowired;
Expand All @@ -34,6 +40,8 @@
@Component
public class LogClientDelegate {

private static final int LOG_CHUNK_SIZE = 8 * 1024 * 1024; // 8 MB

@Autowired
private LocalLogClient localLogClient;
@Autowired
Expand Down Expand Up @@ -66,29 +74,6 @@ public String getPartLogString(TaskInstance taskInstance, int skipLineNum, int l
}
}

/**
* Retrieves the complete log content for a given task instance as a byte array.
* This method first attempts to fetch the log from local storage; if unsuccessful, it tries to obtain the log from remote storage.
*
* @param taskInstance The task instance object, containing information needed for log retrieval.
* @return A byte array containing the complete log content.
*/
public byte[] getWholeLogBytes(TaskInstance taskInstance) {
checkArgs(taskInstance);
if (checkNodeExists(taskInstance)) {
TaskInstanceLogFileDownloadResponse response = localLogClient.getWholeLog(taskInstance);
if (response.getCode() == LogResponseStatus.SUCCESS) {
return response.getLogBytes();
} else {
log.warn("get whole log bytes is not success for task instance {}; reason :{}", taskInstance.getId(),
response.getMessage());
return remoteLogClient.getWholeLog(taskInstance);
}
} else {
return remoteLogClient.getWholeLog(taskInstance);
}
}

private static void checkArgs(TaskInstance taskInstance) {
if (taskInstance == null) {
throw new IllegalArgumentException("canFetchLog task instance is null");
Expand All @@ -109,4 +94,101 @@ private boolean checkNodeExists(TaskInstance taskInstance) {
return exists;
}

/**
* Stream the entire task instance log to {@code outputStream} using bounded chunk RPCs.
*
* <p>Strategy:
* <ul>
* <li>If the worker node is gone, read straight from remote log storage (archive), streamed
* in bounded chunks.</li>
* <li>Otherwise stream via the chunk RPC. If the FIRST chunk fails — e.g. an old worker
* during a rolling upgrade does not implement {@code getTaskInstanceLogFileChunk} —
* fall back to remote log storage; if that also fails, throw an explicit error asking
* for a worker upgrade.</li>
* <li>The legacy whole-file worker RPC is deliberately NEVER used: it reads the entire file
* into the worker's heap before serialization, so a large log can OOM the worker. A
* receiver-side maxFrameSize cannot prevent that allocation, and an old worker (already
* deployed, cannot be patched) offers no way to establish a safe size before the whole
* payload is built — so small-log compatibility with old workers is intentionally not
* preserved either.</li>
* <li>If a failure happens mid-stream (bytes already written), throw IOException to avoid
* corrupting the download.</li>
* </ul>
*/
public void streamWholeLog(final TaskInstance taskInstance, final OutputStream outputStream) throws IOException {
checkArgs(taskInstance);
if (!checkNodeExists(taskInstance)) {
remoteLogClient.streamWholeLog(taskInstance, outputStream);
return;
}
long offset = 0;
while (true) {
final TaskInstanceLogFileDownloadResponse chunk;
try {
chunk = localLogClient.getLogChunk(taskInstance, offset, LOG_CHUNK_SIZE);
} catch (Exception e) {
if (offset > 0) {
throw new IOException("Log streaming failed at offset " + offset, e);
}
log.warn("Chunked log RPC failed for task instance {}, falling back to remote log storage",
taskInstance.getId(), e);
// A MethodInvocationException means the worker ANSWERED but could not dispatch the
// method — that is the old-worker signal (rolling upgrade: the chunk method does
// not exist there), so the upgrade guidance is accurate. Any other transport
// failure (connect refused, timeout) just means the worker is unreachable; blaming
// the worker version would mislead operations.
final String errorMessage =
ExceptionUtils.throwableOfType(e, MethodInvocationException.class) != null
? "Worker upgrade required for large log download: chunked log RPC is not available"
+ " on worker " + taskInstance.getHost() + " and remote log storage also failed"
: "Chunked log RPC to worker " + taskInstance.getHost()
+ " failed (the worker may be down or unreachable)"
+ " and remote log storage also failed";
fallbackToRemoteStorage(taskInstance, outputStream, errorMessage);
return;
}
if (chunk == null || chunk.getCode() != LogResponseStatus.SUCCESS) {
final String failure = chunk == null
? "worker returned no response"
: chunk.getCode() + ": " + chunk.getMessage();
if (offset > 0) {
throw new IOException("Worker chunk failed at offset " + offset + ": " + failure);
}
log.warn("First chunk failed for task instance {} ({}), falling back to remote log storage",
taskInstance.getId(), failure);
// The worker ANSWERED (structured response) — it supports the chunk RPC, so no
// upgrade guidance here; report both failures as they are.
fallbackToRemoteStorage(taskInstance, outputStream,
"Chunked log fetch failed on worker " + taskInstance.getHost() + " for task instance "
+ taskInstance.getId() + " (" + failure + ") and remote log storage also failed");
return;
}
final byte[] data = chunk.getLogBytes();
if (data != null && data.length > 0) {
outputStream.write(data);
offset += data.length;
}
if (chunk.isEof() || (data == null || data.length == 0)) {
return;
}
}
}

/**
* The ONLY fallback of the streaming path: stream the log from remote log storage. If remote
* storage cannot serve the log either, fail with an explicit error — the legacy whole-file
* worker RPC is deliberately never used (it is unbounded on the worker side, see
* {@link #streamWholeLog}). Exceptions from the fallback propagate directly; there is no
* second fallback to re-enter.
*/
private void fallbackToRemoteStorage(final TaskInstance taskInstance,
final OutputStream outputStream,
final String errorMessage) throws IOException {
try {
remoteLogClient.streamWholeLog(taskInstance, outputStream);
} catch (Exception e) {
throw new IOException(errorMessage, e);
}
}

}
Loading
Loading