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
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
import com.here.xyz.jobs.JobClientInfo;
import com.here.xyz.jobs.steps.execution.LambdaBasedStep.LambdaStepRequest.ProcessUpdate;
import com.here.xyz.jobs.steps.execution.StepException;
import com.here.xyz.jobs.steps.execution.db.Database;
import com.here.xyz.jobs.steps.impl.SpaceBasedStep;
import com.here.xyz.jobs.steps.impl.transport.tasks.TaskPayload;
import com.here.xyz.jobs.steps.impl.transport.tasks.TaskProgress;
Expand All @@ -58,10 +59,16 @@
import java.io.IOException;
import java.math.BigDecimal;
import java.math.RoundingMode;
import java.sql.Array;
import java.sql.ResultSet;
import java.sql.SQLException;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
import java.util.Optional;
import java.util.Set;
import java.util.stream.Collectors;

import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;

Expand All @@ -83,7 +90,12 @@
public abstract class TaskedSpaceBasedStep<T extends TaskedSpaceBasedStep, I extends TaskPayload, O extends TaskPayload>
extends SpaceBasedStep<T> {
private static final Logger logger = LogManager.getLogger();
/** Marker output-set key/file prefix used to indicate successful step finalization. */
public static final String FINALIZATION_MARKER = "finalized_marker";
/** Number of consecutive unknown running-query checks before a task is considered retryable. */
public static final Integer MAX_UNKNOWN_TASK_QUERY_CHECKS = 3;
/** Maximum number of retry attempts, for a server side killed single, before failing retry handling. */
public static final Integer MAX_TASK_RETRY_ATTEMPTS = 3;
private TaskedSpaceBasedQueryBuilder taskedSpaceBasedQueryBuilder;

{
Expand Down Expand Up @@ -491,8 +503,9 @@ private void startTask(TaskProgress<I> taskProgressAndItem) throws TooManyResour

prepareTaskQuery(taskProgressAndItem.getTaskId());

runReadQueryAsync(buildTaskQuery(taskProgressAndItem.getTaskId(), (I) taskProgressAndItem.getTaskInput(), failureCallback),
queryRunsOnWriter() ? dbWriter() : dbReader(), 0d/*perItemAcus.doubleValue()*/, false);
runReadQueryAsync(buildTaskQuery(taskProgressAndItem.getTaskId(), (I) taskProgressAndItem.getTaskInput(), failureCallback)
.withLabel(getId() + "#taskId", String.valueOf(taskProgressAndItem.getTaskId())),
queryRunsOnWriter() ? dbWriter() : dbReader(), 0d/*perItemAcus.doubleValue()*/, false);
}
}

Expand Down Expand Up @@ -696,7 +709,7 @@ public AsyncExecutionState getExecutionState() throws UnknownStateException {
return AsyncExecutionState.SUCCEEDED;

try {
TaskProgress taskProgress = getTaskProgress();
TaskProgress<?> taskProgress = getTaskProgress();
return evaluateExecutionState(taskProgress);
}
catch (SQLException e) {
Expand All @@ -709,13 +722,57 @@ public AsyncExecutionState getExecutionState() throws UnknownStateException {
}
}

private AsyncExecutionState evaluateExecutionState(TaskProgress taskProgress) throws UnknownStateException {
private AsyncExecutionState evaluateExecutionState(TaskProgress<?> taskProgress)
throws UnknownStateException, SQLException, WebClientException, TooManyResourcesClaimed {
if (taskProgress.isComplete()) {
//Only log and return RUNNING to avoid double success handling;
infoLog(STEP_ON_STATE_CHECK, "All tasks finalized!");
}if (taskProgress.hasRunningTasks())
// Check if the expected queries are still running.
return super.getExecutionState();
}
if (taskProgress.hasRunningTasks()) {
Set<Integer> expectedTaskIds = taskProgress.getStartedNotFinalizedTaskIds();
Set<String> expectedTaskIdStrings = expectedTaskIds.stream()
.map(String::valueOf)
.collect(Collectors.toSet());

if (expectedTaskIdStrings.isEmpty()) {
infoLog(STEP_ON_STATE_CHECK, "Started tasks are reported but started_not_finalized_task_ids is empty.");
throw new UnknownStateException("No started-not-finalized task ids available for step: " + getGlobalStepId() + " !");
}

boolean runsOnWriter = queryRunsOnWriter();
Database database = runsOnWriter ? dbWriter() : dbReader();
//Check if all expected TaskQueries areRunning. Collect all taskIds where we are not able
//to find a running query.
Set<String> runningTaskIds = SQLQuery.areRunning(
requestResource(database, 0d), database.getRole() != WRITER,
getId() + "#taskId", expectedTaskIdStrings
);

Set<Integer> unknownStateTaskIds = expectedTaskIds.stream()
.filter(taskId -> !runningTaskIds.contains(String.valueOf(taskId)))
.collect(Collectors.toSet());

if (unknownStateTaskIds.isEmpty())
return AsyncExecutionState.RUNNING;

Set<Integer> exceededTaskIds = incrementUnknownQueryStateForTasks(unknownStateTaskIds);
for(int exceededTaskId : exceededTaskIds){
infoLog(STEP_ON_STATE_CHECK, "Unknown queryState of taskId " + exceededTaskId + " has exceeded unknown query state threshold. Retry the task now!");
try {
TaskProgress<I> retryTaskProgress = resetTaskForRetry(exceededTaskId);
if (retryTaskProgress == null) {
throw new StepException("Retry threshold exceeded for taskId " + exceededTaskId + ".");
}
startTask(retryTaskProgress);
} catch (StepException e1){
throw e1;
} catch (Exception e2){
throw new StepException("Not able to retry taskId " + exceededTaskId, e2);
}
}

return AsyncExecutionState.RUNNING;
}
if (taskProgress.hasNoRunningTasks()) {
infoLog(STEP_ON_STATE_CHECK, "No running tasks detected. StartedTasks: " + taskProgress.getStartedTasks() + ","
+ " FinalizedTasks: " + taskProgress.getFinalizedTasks() + " !");
Expand Down Expand Up @@ -783,6 +840,40 @@ private void updateQueryTaskItemOutput(SpaceBasedTaskUpdate update) throws WebCl
, db(WRITER), 0);
}

private Set<Integer> incrementUnknownQueryStateForTasks(Set<Integer> unknownStateTaskIds)
throws WebClientException, SQLException, TooManyResourcesClaimed {
if (unknownStateTaskIds == null || unknownStateTaskIds.isEmpty())
return Set.of();

SQLQuery query = getQueryBuilder().buildIncrementUnknownQueryStateStatement(unknownStateTaskIds, MAX_UNKNOWN_TASK_QUERY_CHECKS);

return runReadQuerySync(query, db(WRITER), 0, rs -> {
Set<Integer> exceededTaskIds = new java.util.LinkedHashSet<>();
while (rs.next())
exceededTaskIds.add(rs.getInt("task_id"));
return exceededTaskIds;
});
}

private TaskProgress<I> resetTaskForRetry(int taskId) throws WebClientException, SQLException, TooManyResourcesClaimed {
return runReadQuerySync(getQueryBuilder().buildResetTaskForRetryStatement(taskId), db(WRITER), 0, rs -> {
if (!rs.next())
throw new SQLException("Task for retry not found: taskId=" + taskId);

int retryAttempt = rs.getInt("retry_attempts");
if (retryAttempt > MAX_TASK_RETRY_ATTEMPTS)
return null;

try {
return new TaskProgress<>(rs.getInt("task_id"),
XyzSerializable.deserialize(rs.getString("task_input"), new TypeReference<I>() {}));
}
catch (JsonProcessingException e) {
throw new StepException("Can not deserialize task_input for retry of taskId=" + taskId + "!", e);
}
});
}

private TaskProgress executeTaskItemQuery(SQLQuery query) throws WebClientException, SQLException, TooManyResourcesClaimed {
TaskProgress taskProgress;
try {
Expand Down Expand Up @@ -816,10 +907,20 @@ private TaskProgress getTaskProgress() throws WebClientException, SQLException,
rs -> {
if (!rs.next())
return null;
return new TaskProgress(rs.getInt("total"), rs.getInt("started"), rs.getInt("finalized"));
return new TaskProgress(rs.getInt("total"), rs.getInt("started"), rs.getInt("finalized"),getIntegerList(rs, "started_not_finalized_task_ids"));
});
}

private Set<Integer> getIntegerList(ResultSet rs, String columnName) throws SQLException {
Array sqlArray = rs.getArray(columnName);

return sqlArray == null
? Set.of()
: Arrays.stream((Object[]) sqlArray.getArray())
.map(value -> ((Number) value).intValue())
.collect(Collectors.toSet());
}

private SQLQuery resetTaskItemWhichAreNotFinalized() {
infoLog(STEP_EXECUTE, "Reset task items for restart.");
return getQueryBuilder().buildResetTaskItemWhichAreNotFinalizedStatement();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,19 +18,27 @@
*/
package com.here.xyz.jobs.steps.impl.transport.tasks;

import java.util.Set;

public class TaskProgress<I> {
private int totalTasks;
private int startedTasks;
private int finalizedTasks;
private Integer taskId;
private I taskInput;
private Set<Integer> startedNotFinalizedTaskIds = Set.of();

public TaskProgress() {}

public TaskProgress(Integer taskId) {
this.taskId = taskId;
}

public TaskProgress(Integer taskId, I taskInput) {
this.taskId = taskId;
this.taskInput = taskInput;
}

public TaskProgress(int totalTasks, int startedTasks, int finalizedTasks) {
this.totalTasks = totalTasks;
this.startedTasks = startedTasks;
Expand All @@ -45,6 +53,13 @@ public TaskProgress(int totalTasks, int startedTasks, int finalizedTasks, Intege
this.taskInput = taskInput;
}

public TaskProgress(int totalTasks, int startedTasks, int finalizedTasks, Set<Integer> startedNotFinalizedTaskIds) {
this.totalTasks = totalTasks;
this.startedTasks = startedTasks;
this.finalizedTasks = finalizedTasks;
this.startedNotFinalizedTaskIds = startedNotFinalizedTaskIds;
}

public int getTotalTasks() {
return totalTasks;
}
Expand Down Expand Up @@ -85,6 +100,14 @@ public void setTaskInput(I taskInput) {
this.taskInput = taskInput;
}

public Set<Integer> getStartedNotFinalizedTaskIds() {
return startedNotFinalizedTaskIds;
}

public void setStartedNotFinalizedTaskIds(Set<Integer> startedNotFinalizedTaskIds) {
this.startedNotFinalizedTaskIds = startedNotFinalizedTaskIds;
}

public boolean isComplete() {
return totalTasks == finalizedTasks;
}
Expand All @@ -109,6 +132,7 @@ public String toString() {
"totalTasks=" + totalTasks +
", startedTasks=" + startedTasks +
", finalizedTasks=" + finalizedTasks +
", startedNotFinalizedTaskIds=" + startedNotFinalizedTaskIds +
", taskId=" + taskId +
", taskInput=" + taskInput +
'}';
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@
import com.here.xyz.models.hub.Space;
import com.here.xyz.util.db.SQLQuery;

import java.util.Set;

import static com.here.xyz.events.ContextAwareEvent.SpaceContext;
import static com.here.xyz.jobs.steps.impl.transport.TaskedSpaceBasedStep.SpaceBasedTaskUpdate;

Expand All @@ -42,10 +44,13 @@ public SQLQuery buildTaskTableStatement() {
task_output JSONB,
started BOOLEAN DEFAULT false,
finalized BOOLEAN DEFAULT false,
unknown_query_state_occurrences INTEGER DEFAULT 0,
retry_attempts INTEGER DEFAULT 0,
started_at TIMESTAMP DEFAULT NULL,
updated_at TIMESTAMP DEFAULT NULL,
CONSTRAINT ${primaryKey} PRIMARY KEY (task_id)
);
""")
//TODO: CHECK CONSTRAINT!!
.withVariable("table", getTemporaryJobTableName())
.withVariable("schema", schema)
.withVariable("primaryKey", getTemporaryJobTableName() + "_primKey");
Expand All @@ -54,15 +59,17 @@ public SQLQuery buildTaskTableStatement() {
public SQLQuery buildUpdateTaskItemOutputStatement(SpaceBasedTaskUpdate update) {
return new SQLQuery("""
UPDATE ${schema}.${table}
SET task_output = (
COALESCE(task_output, '{}'::JSONB) || #{taskUpdate}::JSONB
) || jsonb_build_object(
'taskOutput',
COALESCE(task_output->'taskOutput', '{}'::JSONB)
|| COALESCE((#{taskUpdate}::JSONB)->'taskOutput', '{}'::JSONB)
)
WHERE task_id = #{taskId};
""")
SET
updated_at = now(),
task_output = (
COALESCE(task_output, '{}'::JSONB) || #{taskUpdate}::JSONB
) || jsonb_build_object(
'taskOutput',
COALESCE(task_output->'taskOutput', '{}'::JSONB)
|| COALESCE((#{taskUpdate}::JSONB)->'taskOutput', '{}'::JSONB)
)
WHERE task_id = #{taskId};
""")
.withVariable("schema", schema)
.withVariable("table", getTemporaryJobTableName())
.withNamedParameter("taskId", update.taskId)
Expand All @@ -72,7 +79,8 @@ public SQLQuery buildUpdateTaskItemOutputStatement(SpaceBasedTaskUpdate update)
public SQLQuery buildResetTaskItemWhichAreNotFinalizedStatement() {
return new SQLQuery("""
UPDATE ${schema}.${table} t
SET started = false
SET started = false,
updated_at = now()
WHERE started = true AND finalized = false;
""")
.withVariable("schema", schema)
Expand All @@ -89,7 +97,8 @@ public SQLQuery retrieveTaskStatisticsQuery() {
return new SQLQuery("""
SELECT COUNT(1) as total,
SUM((started = true)::int) as started,
SUM((finalized = true)::int) as finalized
SUM((finalized = true)::int) as finalized,
ARRAY_AGG(task_id) FILTER (WHERE started = true AND finalized = false) as started_not_finalized_task_ids
FROM ${schema}.${table};
""")
.withVariable("schema", schema)
Expand Down Expand Up @@ -144,6 +153,42 @@ public SQLQuery buildTemporaryJobTableDropStatement() {
.withVariable("schema", schema);
}

public SQLQuery buildIncrementUnknownQueryStateStatement(Set<Integer> unknownStateTaskIds, int maxUnknownTaskQueryChecks) {
return new SQLQuery("""
WITH updated AS (
UPDATE ${schema}.${table}
SET unknown_query_state_occurrences = COALESCE(unknown_query_state_occurrences, 0) + 1,
updated_at = now()
WHERE task_id IN (
SELECT value::INT
FROM jsonb_array_elements_text(#{missingTaskIds}::JSONB)
)
RETURNING task_id, unknown_query_state_occurrences
)
SELECT task_id
FROM updated
WHERE unknown_query_state_occurrences >= #{maxUnknownTaskQueryChecks};
""")
.withVariable("schema", schema)
.withVariable("table", getTemporaryJobTableName())
.withNamedParameter("missingTaskIds", XyzSerializable.serialize(unknownStateTaskIds))
.withNamedParameter("maxUnknownTaskQueryChecks", maxUnknownTaskQueryChecks);
}

public SQLQuery buildResetTaskForRetryStatement(int taskId) {
return new SQLQuery("""
UPDATE ${schema}.${table}
SET unknown_query_state_occurrences = 0,
updated_at = now(),
retry_attempts = retry_attempts + 1
WHERE task_id = #{taskId}
RETURNING task_id, task_input, retry_attempts;
""")
.withVariable("schema", schema)
.withVariable("table", getTemporaryJobTableName())
.withNamedParameter("taskId", taskId);
}

private String getTemporaryJobTableName(String stepId) {
return JOB_DATA_PREFIX + stepId;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -569,7 +569,7 @@ BEGIN
IF task_item.task_id IS NOT NULL THEN
EXECUTE format(
'UPDATE %1$s C
SET started = true
SET started = true, started_at = now()
WHERE C.task_id = %2$L;',
get_table_reference(ctx->>'schema', ctx->>'stepId' ,'JOB_TABLE'),
task_item.task_id
Expand Down Expand Up @@ -663,7 +663,7 @@ BEGIN
COALESCE(t.task_output->''taskOutput'', ''{}''::JSONB)
|| COALESCE((%2$L::JSONB)->''taskOutput'', ''{}''::JSONB)
),
finalized = %3$L
updated_at = now(), finalized = %3$L
WHERE task_id = %4$L;',
get_table_reference(ctx->>'schema', ctx->>'stepId', 'JOB_TABLE'),
p_task_output::TEXT,
Expand Down
Loading
Loading