diff --git a/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml b/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml index 94347bf9e0f2..433759424509 100644 --- a/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml +++ b/.github/workflows/beam_PostCommit_Java_Delta_IO_Dataflow.yml @@ -64,8 +64,6 @@ jobs: job_phrase: ["Run PostCommit Java Delta IO Dataflow"] steps: - uses: actions/checkout@v7 - with: - persist-credentials: false - name: Setup repository uses: ./.github/actions/setup-action with: diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java index af8e7dd50c95..810e9d20ed77 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java @@ -135,13 +135,13 @@ public Optional getWorkItem() throws IOException { final String stage; if (work.getMapTask() != null) { - stage = work.getMapTask().getStageName(); + stage = work.getMapTask().getSystemName(); logger.info("Starting MapTask stage {}", stage); } else if (work.getSeqMapTask() != null) { - stage = work.getSeqMapTask().getStageName(); + stage = work.getSeqMapTask().getSystemName(); logger.info("Starting SeqMapTask stage {}", stage); } else if (work.getSourceOperationTask() != null) { - stage = work.getSourceOperationTask().getStageName(); + stage = work.getSourceOperationTask().getSystemName(); logger.info("Starting SourceOperationTask stage {}", stage); } else { stage = null; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java index 2339430464c7..995ab7d778a6 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java @@ -924,7 +924,7 @@ static StreamingDataflowWorker forTesting( mapTask, workExecutor, stateNameMap, - stateCache.forComputation(mapTask.getStageName()))); + stateCache.forComputation(mapTask.getStageName(), mapTask.getSystemName()))); MemoryMonitor memoryMonitor = MemoryMonitor.fromOptions(options); FailureTracker failureTracker = options.isEnableStreamingEngine() @@ -1197,26 +1197,23 @@ void stop() { } private void onCompleteCommit(CompleteCommit completeCommit) { + Optional computationState = + computationStateCache.getIfPresent(completeCommit.computationId()); if (completeCommit.status() != Windmill.CommitStatus.OK) { readerCache.invalidateReader( WindmillComputationKey.create( completeCommit.computationId(), completeCommit.shardedKey())); - stateCache - .forComputation(completeCommit.computationId()) - .invalidate(completeCommit.shardedKey()); + computationState.ifPresent( + state -> + stateCache + .forComputation(completeCommit.computationId(), state.getSystemName()) + .invalidate(completeCommit.shardedKey())); } - computationStateCache - .getIfPresent(completeCommit.computationId()) - .ifPresent( - state -> { - if (completeCommit.retryableFailure()) { - state.reexecuteActiveWork(completeCommit.shardedKey(), completeCommit.workId()); - } else { - state.completeWorkAndScheduleNextWorkForKey( - completeCommit.shardedKey(), completeCommit.workId()); - } - }); + computationState.ifPresent( + state -> + state.completeWorkAndScheduleNextWorkForKey( + completeCommit.shardedKey(), completeCommit.workId())); } @AutoValue diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java index 6894ac20ef97..74b17b52bcf3 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java @@ -52,7 +52,6 @@ import org.apache.beam.runners.dataflow.worker.counters.NameContext; import org.apache.beam.runners.dataflow.worker.profiler.ScopedProfiler.ProfileScope; import org.apache.beam.runners.dataflow.worker.streaming.BoundedQueueExecutorWorkHandle; -import org.apache.beam.runners.dataflow.worker.streaming.ExecutableWork; import org.apache.beam.runners.dataflow.worker.streaming.KeyCommitTooLargeException; import org.apache.beam.runners.dataflow.worker.streaming.Watermarks; import org.apache.beam.runners.dataflow.worker.streaming.Work; @@ -268,6 +267,10 @@ public final long getBacklogBytes() { return backlogBytes; } + public String getSystemName() { + return systemName; + } + public long getMaxOutputKeyBytes() { return operationalLimits.getMaxOutputKeyBytes(); } @@ -585,7 +588,7 @@ public void invalidateCache() { } catch (IOException e) { Windmill.WorkItem workItem = getWorkItem(); long shardingKey = workItem != null ? workItem.getShardingKey() : -1L; - LOG.warn("Failed to close reader for {}-{}", computationId, shardingKey, e); + LOG.warn("Failed to close reader for {}-{}", systemName, shardingKey, e); } } activeReader = null; @@ -726,11 +729,6 @@ private void validateCommitRequestSize() { buildWorkItemTruncationRequestBuilder(currentWork, estimatedCommitSize); currentBuilder.clear(); currentBuilder.mergeFrom(truncationBuilder.build()); - - // TODO: throw and retry when truncation is not on a single key bundle. - checkState( - !multiKeyBundleOptions.multiKeyBundleEnabled(), - "Commit truncation not implemented for multikey bundles"); } private Windmill.WorkItemCommitRequest.Builder buildWorkItemTruncationRequestBuilder( diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java index abe5f96bb7f4..178509594660 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java @@ -265,20 +265,26 @@ public long add(WindowedValue data) throws IOException { } if (key.size() > context.getMaxOutputKeyBytes()) { if (context.throwExceptionsForLargeOutput()) { - throw new OutputTooLargeException("Key too large: " + key.size()); + throw new OutputTooLargeException( + String.format( + "Key for system %s too large: %s", context.getSystemName(), key.size())); } else { LOG.error( - "Trying to output too large key with size {}. Limit is {}. See https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception. Running with --experiments=throw_exceptions_on_large_output will instead throw an OutputTooLargeException which may be caught in user code.", + "Trying to output too large key for system {} with size {}. Limit is {}. See https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception. Running with --experiments=throw_exceptions_on_large_output will instead throw an OutputTooLargeException which may be caught in user code.", + context.getSystemName(), key.size(), context.getMaxOutputKeyBytes()); } } if (value.size() > context.getMaxOutputValueBytes()) { if (context.throwExceptionsForLargeOutput()) { - throw new OutputTooLargeException("Value too large: " + value.size()); + throw new OutputTooLargeException( + String.format( + "Value for system %s too large: %s", context.getSystemName(), value.size())); } else { LOG.error( - "Trying to output too large value with size {}. Limit is {}. See https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception. Running with --experiments=throw_exceptions_on_large_output will instead throw an OutputTooLargeException which may be caught in user code.", + "Trying to output too large value for system {} with size {}. Limit is {}. See https://cloud.google.com/dataflow/docs/guides/common-errors#key-commit-too-large-exception. Running with --experiments=throw_exceptions_on_large_output will instead throw an OutputTooLargeException which may be caught in user code.", + context.getSystemName(), value.size(), context.getMaxOutputValueBytes()); } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java index de4082581293..f0150cf73eb3 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ActiveWorkState.java @@ -185,7 +185,7 @@ synchronized void failWorkForKey(ImmutableList failedWork executableWork.work().setFailed(); LOG.debug( "Failing work {} {}. The work will be retried and is not lost.", - computationStateCache.getComputation(), + computationStateCache.getSystemName(), failedId); } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java index 8020eda1b25d..5e850d4312ea 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java @@ -69,6 +69,10 @@ public String getComputationId() { return computationId; } + public String getSystemName() { + return mapTask.getSystemName(); + } + public MapTask getMapTask() { return mapTask; } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java index 4b4acb73f4a7..e6f902a65bdc 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java @@ -28,6 +28,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutionException; +import java.util.function.BiFunction; import java.util.function.Function; import javax.annotation.concurrent.ThreadSafe; import org.apache.beam.runners.dataflow.worker.apiary.FixMultiOutputInfosOnParDoInstructions; @@ -77,7 +78,8 @@ private ComputationStateCache( public static ComputationStateCache create( ComputationConfig.Fetcher computationConfigFetcher, BoundedQueueExecutor workUnitExecutor, - Function perComputationStateCacheViewFactory, + BiFunction + perComputationStateCacheViewFactory, IdGenerator idGenerator) { Function fixMultiOutputInfosOnParDoInstructions = new FixMultiOutputInfosOnParDoInstructions(idGenerator); @@ -105,7 +107,8 @@ public ComputationState load(String computationId) { fixMultiOutputInfosOnParDoInstructions.apply(computationConfig.mapTask()), workUnitExecutor, transformUserNameToStateFamilyForComputation, - perComputationStateCacheViewFactory.apply(computationId)); + perComputationStateCacheViewFactory.apply( + computationId, computationConfig.mapTask().getSystemName())); } }), fixMultiOutputInfosOnParDoInstructions, @@ -116,7 +119,8 @@ public ComputationState load(String computationId) { public static ComputationStateCache forTesting( ComputationConfig.Fetcher computationConfigFetcher, BoundedQueueExecutor workUnitExecutor, - Function perComputationStateCacheViewFactory, + BiFunction + perComputationStateCacheViewFactory, IdGenerator idGenerator, ConcurrentMap pipelineUserNameToStateFamilyNameMap) { ComputationStateCache cache = @@ -205,7 +209,7 @@ public void closeAndInvalidateAll() { public void appendSummaryHtml(PrintWriter writer) { writer.println("

Specs

"); for (ComputationState computationState : getAllPresentComputations()) { - writer.println("

" + computationState.getComputationId() + "

"); + writer.println("

" + computationState.getSystemName() + "

"); writer.print(""); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java index 901e2d235f85..0580b7a0b05b 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java @@ -59,7 +59,7 @@ public void appendSummaryHtml(PrintWriter writer) { writer.println("Active Keys:
"); for (ComputationState computationState : allComputationStates.get()) { - writer.print(computationState.getComputationId()); + writer.print(computationState.getSystemName()); writer.print(":
"); computationState.printActiveWork(writer); writer.println("
"); diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java index bbd6cfc9432b..aba9835b9e70 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java @@ -66,9 +66,11 @@ public final String computationId() { return computationState().getComputationId(); } - public @Nullable WorkItemCommitRequest singleKeyRequest() { - return singleKeyRequest; - }; + public final String systemName() { + return computationState().getSystemName(); + } + + public abstract WorkItemCommitRequest request(); public ComputationState computationState() { return computationState; diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java index 8ac9b1593c54..400b4027f184 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java @@ -112,7 +112,12 @@ public void commit(Commit commit) { // Do this check after adding to commitQueue, else commitQueue.put() can race with // drainCommitQueue() in stop() and leave commits orphaned in the queue. if (!this.isRunning.get()) { - LOG.debug("Trying to queue commit on shutdown, failing commit={}", commit); + LOG.debug( + "Trying to queue commit on shutdown, failing commit=[systemName={}, shardingKey={}," + + " workId={} ].", + commit.systemName(), + commit.work().getShardedKey(), + commit.work().id()); drainCommitQueue(); } } diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java index 7515db000852..ff62d12a8fa2 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java @@ -170,8 +170,8 @@ public CacheStats getCacheStats() { } /** Returns a per-computation view of the state cache. */ - public ForComputation forComputation(String computation) { - return new ForComputation(computation); + public ForComputation forComputation(String computation, String systemName) { + return new ForComputation(computation, systemName); } /** Print summary statistics of the cache to the given {@link PrintWriter}. */ @@ -353,9 +353,11 @@ private Optional value() { public class ForComputation { private final String computation; + private final String systemName; - private ForComputation(String computation) { + private ForComputation(String computation, String systemName) { this.computation = computation; + this.systemName = systemName; } /** Returns the computation associated to this class. */ @@ -363,6 +365,11 @@ public String getComputation() { return this.computation; } + /** Returns the system name associated to this class. */ + public String getSystemName() { + return this.systemName; + } + /** Invalidate all cache entries for this computation and {@code processingKey}. */ public void invalidate(ByteString processingKey, long shardingKey) { WindmillComputationKey key = diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java index b51512252e37..f0e5ab019420 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java @@ -20,6 +20,7 @@ import static org.apache.beam.runners.dataflow.DataflowRunner.hasExperiment; import com.google.api.services.dataflow.model.MapTask; +import java.util.function.BiFunction; import java.util.function.Function; import org.apache.beam.runners.dataflow.internal.CustomSources; import org.apache.beam.runners.dataflow.options.DataflowWorkerHarnessOptions; @@ -82,12 +83,11 @@ final class ComputationWorkExecutorFactory { private final DataflowWorkerHarnessOptions options; private final DataflowMapTaskExecutorFactory mapTaskExecutorFactory; private final ReaderCache readerCache; - private final Function stateCacheFactory; + private final BiFunction stateCacheFactory; private final ReaderRegistry readerRegistry; private final SinkRegistry sinkRegistry; private final DataflowExecutionStateSampler sampler; private final CounterSet pendingDeltaCounters; - private final SideInputStateFetcherFactory sideInputStateFetcherFactory; private final StreamingCounters streamingCounters; private final FailureTracker failureTracker; @@ -112,7 +112,7 @@ final class ComputationWorkExecutorFactory { DataflowWorkerHarnessOptions options, DataflowMapTaskExecutorFactory mapTaskExecutorFactory, ReaderCache readerCache, - Function stateCacheFactory, + BiFunction stateCacheFactory, DataflowExecutionStateSampler sampler, StreamingCounters streamingCounters, FailureTracker failureTracker, @@ -287,7 +287,7 @@ private StreamingModeExecutionContext createExecutionContext( computationId, readerCache, computationState.getTransformUserNameToStateFamily(), - stateCacheFactory.apply(computationId), + stateCacheFactory.apply(computationId, stageInfo.systemName()), stageInfo.metricsContainerRegistry(), executionStateTracker, stageInfo.executionStateRegistry(), diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java index 9e8265e509af..2ac3ebb706a9 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java @@ -21,13 +21,12 @@ import com.google.api.services.dataflow.model.MapTask; import com.google.auto.value.AutoValue; -import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; -import java.util.function.Function; +import java.util.function.BiFunction; import java.util.function.Supplier; import javax.annotation.concurrent.ThreadSafe; import org.apache.beam.repackaged.core.org.apache.commons.lang3.tuple.Pair; @@ -88,7 +87,6 @@ public class StreamingWorkScheduler { private final ConcurrentMap stageInfoMap; private final DataflowExecutionStateSampler sampler; private final BoundedQueueExecutor workExecutor; - private final MultiKeyBundleOptions multiKeyBundleOptions; public StreamingWorkScheduler( Supplier clock, @@ -98,8 +96,7 @@ public StreamingWorkScheduler( StreamingCommitFinalizer commitFinalizer, StreamingCounters streamingCounters, ConcurrentMap stageInfoMap, - DataflowExecutionStateSampler sampler, - MultiKeyBundleOptions multiKeyBundleOptions) { + DataflowExecutionStateSampler sampler) { this.clock = clock; this.workExecutor = workExecutor; this.computationWorkExecutorFactory = computationWorkExecutorFactory; @@ -108,7 +105,6 @@ public StreamingWorkScheduler( this.streamingCounters = streamingCounters; this.stageInfoMap = stageInfoMap; this.sampler = sampler; - this.multiKeyBundleOptions = multiKeyBundleOptions; } public static StreamingWorkScheduler create( @@ -119,7 +115,7 @@ public static StreamingWorkScheduler create( DataflowMapTaskExecutorFactory mapTaskExecutorFactory, BoundedQueueExecutor workExecutor, ScheduledExecutorService commitFinalizerCleanupExecutor, - Function stateCacheFactory, + BiFunction stateCacheFactory, FailureTracker failureTracker, WorkFailureProcessor workFailureProcessor, StreamingCounters streamingCounters, @@ -154,8 +150,7 @@ public static StreamingWorkScheduler create( StreamingCommitFinalizer.create(workExecutor, commitFinalizerCleanupExecutor), streamingCounters, stageInfoMap, - sampler, - multiKeyBundleOptions); + sampler); } private static long computeShuffleBytesRead(Windmill.WorkItem workItem) { @@ -175,8 +170,14 @@ private static Windmill.WorkItemCommitRequest.Builder initializeOutputBuilder( .setCacheToken(workItem.getCacheToken()); } - private static void setLoggingContextComputation(@Nullable String computationId) { - DataflowWorkerLoggingMDC.setStageName(computationId); + /** Sets the stage name and workId of the Thread executing the {@link Work} for logging. */ + private static void setUpWorkLoggingContext(String workLatencyTrackingId, String systemName) { + setLoggingContextWorkId(workLatencyTrackingId); + setLoggingContextSystemName(systemName); + } + + private static void setLoggingContextSystemName(@Nullable String systemName) { + DataflowWorkerLoggingMDC.setStageName(systemName); } private static void setLoggingContextWorkId(@Nullable String workLatencyTrackingId) { @@ -226,15 +227,11 @@ public void queueAppliedFinalizeIds(ImmutableList appliedFinalizeIds) { private void processWork( ComputationState computationState, Work work, BoundedQueueExecutorWorkHandle handle) { Windmill.WorkItem workItem = work.getWorkItem(); - String computationId = computationState.getComputationId(); - LOG.debug("Starting processing for {}:\n{}", computationId, work); - setLoggingContextComputation(computationId); - KeyTransitionListener keyTransitionListener = createKeyTransitionListener(); - keyTransitionListener.onKeyTransition(null, work); - - // Before any processing starts, call any pending OnCommit callbacks. Nothing that requires - // cleanup should be done before this, since we might exit early here. - commitFinalizer.finalizeCommits(workItem.getSourceState().getFinalizeIdsList()); + String systemName = computationState.getSystemName(); + work.setProcessingThreadName(Thread.currentThread().getName()); + work.setState(Work.State.PROCESSING); + setUpWorkLoggingContext(work.getLatencyTrackingId(), systemName); + LOG.debug("Starting processing for {}:\n{}", systemName, work); if (workItem.getSourceState().getOnlyFinalize()) { handleOnlyFinalize(computationState, work, workItem); @@ -263,7 +260,21 @@ private void processWork( recordProcessingStats(workBatch, workItemCommits, executeWorkResult.stateBytesRead()); LOG.debug("Processing done for work batch size: {}", workBatch.size()); } catch (Throwable t) { - handleProcessWorkFailure(computationState, handle.getWorkBatch(), computationId, work, t); + // OutOfMemoryError that are caught will be rethrown and trigger jvm termination. + try { + workFailureProcessor.logAndProcessFailure( + systemName, + ExecutableWork.create(work, (retry, h) -> processWork(computationState, retry, h)), + t, + invalidWork -> + computationState.completeWorkAndScheduleNextWorkForKey( + invalidWork.getShardedKey(), invalidWork.id())); + } catch (OutOfMemoryError oom) { + throw oom; + } catch (Throwable t2) { + LOG.warn("Failed to process work failure safely for work {}", work.id(), t2); + throw ExceptionUtils.safeWrapThrowableAsException(t2); + } } finally { List processedWorkBatch = workBatch != null ? workBatch : ImmutableList.of(work); // Update total processing time counters. Updating in finally clause ensures that @@ -271,7 +282,7 @@ private void processWork( recordProcessingTime(stageInfo, processedWorkBatch, processingStartTimeNanos); setLoggingContextWorkId(null); - setLoggingContextComputation(null); + setLoggingContextSystemName(null); sampler.resetForWorkId(work.getLatencyTrackingId()); for (Work w : processedWorkBatch) { w.setProcessingThreadName(""); @@ -453,34 +464,6 @@ private void commitSingleKeyWork( work.queueCommit(commitRequestWithAttributions, computationState); } - private void handleProcessWorkFailure( - ComputationState computationState, - List failedBatch, - String computationId, - Work primaryWork, - Throwable t) { - try { - List executableWorks = new ArrayList<>(); - for (Work w : failedBatch) { - executableWorks.add( - ExecutableWork.create(w, (retry, h) -> processWork(computationState, retry, h))); - } - - workFailureProcessor.logAndProcessFailureBatch( - computationId, - executableWorks, - t, - invalidWork -> - computationState.completeWorkAndScheduleNextWorkForKey( - invalidWork.getShardedKey(), invalidWork.id())); - } catch (OutOfMemoryError oom) { - throw oom; - } catch (Throwable t2) { - LOG.warn("Failed to process work failure safely for work {}", primaryWork.id(), t2); - throw ExceptionUtils.safeWrapThrowableAsException(t2); - } - } - private void recordProcessingTime( StageInfo stageInfo, List workBatch, long processingStartTimeNanos) { long processingTimeMsecs = diff --git a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessor.java b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessor.java index 8af1840faf92..c9c44386c187 100644 --- a/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessor.java +++ b/runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessor.java @@ -98,40 +98,25 @@ private static boolean isOutOfMemoryError(@Nullable Throwable t) { return false; } - public void logAndProcessFailureBatch( - String computationId, - List executableWorks, - Throwable t, - Consumer onInvalidWork) + /** + * Processes failures caused by thrown exceptions that occur during execution of {@link Work}. May + * attempt to retry execution of the {@link Work} or drop it if it is invalid. + */ + public void logAndProcessFailure( + String systemName, ExecutableWork executableWork, Throwable t, Consumer onInvalidWork) throws Throwable { - List worksToRetryLocally = new java.util.ArrayList<>(); - - for (ExecutableWork executableWork : executableWorks) { - switch (evaluateRetry(computationId, executableWork.work(), t)) { - case DO_NOT_RETRY: - // Consider the item invalid. It will eventually be retried by Windmill if it still needs - // to be processed. - onInvalidWork.accept(executableWork.work()); - break; - case RETRY_LOCALLY: - // Try again after some delay and at the end of the queue to avoid a tight loop. - worksToRetryLocally.add(executableWork); - break; - case RETHROW_THROWABLE: - throw t; - } - } - - executeWithDelay(worksToRetryLocally); - } - - private void executeWithDelay(List worksToRetryLocally) { - if (!worksToRetryLocally.isEmpty()) { - // Sleep ONCE for the entire batch delay to avoid sequential thread blocks - Uninterruptibles.sleepUninterruptibly(retryLocallyDelayMs, TimeUnit.MILLISECONDS); - for (ExecutableWork ew : worksToRetryLocally) { - workUnitExecutor.forceExecute(ew, ew.work().getSerializedWorkItemSize()); - } + switch (evaluateRetry(systemName, executableWork.work(), t)) { + case DO_NOT_RETRY: + // Consider the item invalid. It will eventually be retried by Windmill if it still needs to + // be processed. + onInvalidWork.accept(executableWork.work()); + break; + case RETRY_LOCALLY: + // Try again after some delay and at the end of the queue to avoid a tight loop. + executeWithDelay(retryLocallyDelayMs, executableWork); + break; + case RETHROW_THROWABLE: + throw t; } } @@ -148,12 +133,22 @@ private enum RetryEvaluation { RETHROW_THROWABLE, } - private RetryEvaluation evaluateRetry(String computationId, Work work, Throwable t) { - if (work.isFailed()) { + private RetryEvaluation evaluateRetry(String systemName, Work work, Throwable t) { + @Nullable final Throwable cause = t.getCause(); + Throwable parsedException = (t instanceof UserCodeException && cause != null) ? cause : t; + if (KeyTokenInvalidException.isKeyTokenInvalidException(parsedException)) { + LOG.debug( + "Execution of work for system '{}' on sharding key '{}' failed due to token expiration. " + + "Work will not be retried locally.", + systemName, + work.getWorkItem().getShardingKey()); + return RetryEvaluation.DO_NOT_RETRY; + } + if (WorkItemCancelledException.isWorkItemCancelledException(parsedException)) { LOG.debug( - "Execution of work for computation '{}' on sharding key '{}' failed. " - + "Work is already marked as failed, not retrying locally.", - computationId, + "Execution of work for system '{}' on sharding key '{}' failed. " + + "Work will not be retried locally.", + systemName, work.getWorkItem().getShardingKey()); return RetryEvaluation.DO_NOT_RETRY; } @@ -166,30 +161,30 @@ private RetryEvaluation evaluateRetry(String computationId, Work work, Throwable if (isOutOfMemoryError(parsedException)) { String heapDump = tryToDumpHeap(); LOG.error( - "Execution of work for computation '{}' for sharding key '{}' failed with out-of-memory. " + "Execution of work for system '{}' for sharding key '{}' failed with out-of-memory. " + "Work will not be retried locally. Heap dump {}.", - computationId, + systemName, work.getWorkItem().getShardingKey(), heapDump, parsedException); return RetryEvaluation.RETHROW_THROWABLE; } - if (!failureTracker.trackFailure(computationId, work.getWorkItem(), parsedException)) { + if (!failureTracker.trackFailure(systemName, work.getWorkItem(), parsedException)) { LOG.error( - "Execution of work for computation '{}' on sharding key '{}' failed with uncaught exception, " + "Execution of work for system '{}' on sharding key '{}' failed with uncaught exception, " + "and Windmill indicated not to retry locally.", - computationId, + systemName, work.getWorkItem().getShardingKey(), parsedException); return RetryEvaluation.DO_NOT_RETRY; } if (elapsedTimeSinceStart.isLongerThan(MAX_LOCAL_PROCESSING_RETRY_DURATION)) { LOG.error( - "Execution of work for computation '{}' for sharding key '{}' failed with uncaught exception, " + "Execution of work for system '{}' for sharding key '{}' failed with uncaught exception, " + "and it will not be retried locally because the elapsed time since start {} " + "exceeds {}.", - computationId, + systemName, work.getWorkItem().getShardingKey(), elapsedTimeSinceStart, MAX_LOCAL_PROCESSING_RETRY_DURATION, @@ -197,9 +192,9 @@ private RetryEvaluation evaluateRetry(String computationId, Work work, Throwable return RetryEvaluation.DO_NOT_RETRY; } LOG.error( - "Execution of work for computation '{}' on sharding key '{}' failed with uncaught exception. " + "Execution of work for system '{}' on sharding key '{}' failed with uncaught exception. " + "Work will be retried locally.", - computationId, + systemName, work.getWorkItem().getShardingKey(), parsedException); return RetryEvaluation.RETRY_LOCALLY; diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java index 9ed705550bc6..1453c438c2c9 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorkerTest.java @@ -1354,196 +1354,6 @@ public void testMultiKeyCommit_success() throws Exception { worker.stop(); } - @Test - public void testMultiKeyCommit_elementFailure() throws Exception { - if (!streamingEngine) { - return; - } - StreamingDataflowWorker worker = makeMultiKeyEnabledWorker(); - worker.start(); - - String batchInputText = - "work {" - + " computation_id: \"" - + DEFAULT_COMPUTATION_ID - + "\"" - + " input_data_watermark: 0" - + " work {" - + " key: \"key1\"" - + " sharding_key: 1" - + " work_token: 1" - + " cache_token: 2" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data1\"" - + " }" - + " }" - + " }" - + " work {" - + " key: \"key2\"" - + " sharding_key: 2" - + " work_token: 2" - + " cache_token: 3" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data2\"" - + " }" - + " }" - + " }" - + " work {" - + " key: \"key3\"" - + " sharding_key: 3" - + " work_token: 3" - + " cache_token: 4" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data3\"" - + " }" - + " }" - + " }" - + "}"; - Windmill.GetWorkResponse batchInput = - buildInput( - batchInputText, - CoderUtils.encodeToByteArray( - CollectionCoder.of(IntervalWindow.getCoder()), - Collections.singletonList(DEFAULT_WINDOW))); - - server - .whenGetDataCalled() - .answerByDefault( - StreamingDataflowWorkerTest.emptyDataResponderWithFailedWorkTokens(Set.of(2L))); - - server.whenGetWorkCalled().thenReturn(batchInput); - - Map result = server.waitForAndGetCommits(2); - - assertTrue(result.containsKey(1L)); - assertTrue(result.containsKey(3L)); - assertFalse(result.containsKey(2L)); - - List multiKeyCommits = - server.getMultiKeyCommitsReceived(); - assertEquals(1, multiKeyCommits.size()); - Windmill.MultiKeyWorkItemCommitRequest multiKeyCommit = multiKeyCommits.get(0); - assertEquals(2, multiKeyCommit.getRequestsCount()); - assertEquals(3, multiKeyCommit.getRequests(0).getWorkToken()); - assertEquals(1, multiKeyCommit.getRequests(1).getWorkToken()); - - worker.stop(); - } - - @Test - public void testCompleteCommit_retryableFailureTriggersReExecution() throws Exception { - if (!streamingEngine) { - return; - } - StreamingDataflowWorker worker = makeMultiKeyEnabledWorker(); - worker.start(); - - String batchInputText = - "work {" - + " computation_id: \"" - + DEFAULT_COMPUTATION_ID - + "\"" - + " input_data_watermark: 0" - + " work {" - + " key: \"key1\"" - + " sharding_key: 1" - + " work_token: 1" - + " cache_token: 2" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data1\"" - + " }" - + " }" - + " }" - + " work {" - + " key: \"key2\"" - + " sharding_key: 2" - + " work_token: 2" - + " cache_token: 3" - + " key_group { high: 0 low: 1 }" - + " message_bundles {" - + " source_computation_id: \"" - + DEFAULT_SOURCE_COMPUTATION_ID - + "\"" - + " messages {" - + " timestamp: 0" - + " data: \"data2\"" - + " }" - + " }" - + " }" - + "}"; - Windmill.GetWorkResponse batchInput = - buildInput( - batchInputText, - CoderUtils.encodeToByteArray( - CollectionCoder.of(IntervalWindow.getCoder()), - Collections.singletonList(DEFAULT_WINDOW))); - - server - .whenGetDataCalled() - .answerByDefault( - StreamingDataflowWorkerTest.emptyDataResponderWithFailedWorkTokens(Set.of(2L))); - - server.whenGetWorkCalled().thenReturn(batchInput); - - Map result = server.waitForAndGetCommits(1); - - assertTrue(result.containsKey(1L)); - assertFalse(result.containsKey(2L)); - - List multiKeyCommits = - server.getMultiKeyCommitsReceived(); - assertEquals(1, multiKeyCommits.size()); - Windmill.MultiKeyWorkItemCommitRequest multiKeyCommit = multiKeyCommits.get(0); - assertEquals(1, multiKeyCommit.getRequestsCount()); - assertEquals(1, multiKeyCommit.getRequests(0).getWorkToken()); - - worker.stop(); - } - - private StreamingDataflowWorker makeMultiKeyEnabledWorker() { - KvCoder kvCoder = KvCoder.of(StringUtf8Coder.of(), StringUtf8Coder.of()); - - List instructions = - Arrays.asList( - makeSourceInstruction(kvCoder), - makeDoFnInstruction(new WorkDoFn(), 0, kvCoder), - makeSinkInstruction(kvCoder, 1)); - - StreamingDataflowWorker worker = - makeWorker( - defaultWorkerParams( - "--experiments=unstable_enable_multi_key_bundle,windmill_max_key_group_batch_time_ms=50000", - "--numberOfWorkerHarnessThreads=1") - .setLocalRetryTimeoutMs(100) - .setInstructions(instructions) - .build()); - return worker; - } - private void runKeyCommitTooLargeExceptionTest( StreamingDataflowWorkerTestParams.Builder workerParams, boolean expectKeyInErrorMessage) throws Exception { diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java index c5efcea4e47c..6f10f6e3749f 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContextTest.java @@ -134,7 +134,7 @@ private StreamingModeExecutionContext createExecutionContext( WindmillStateCache.builder() .setSizeMb(options.getWorkerCacheMb()) .build() - .forComputation("comp"), + .forComputation("comp", "systemName"), StreamingStepMetricsContainer.createRegistry(), new DataflowExecutionStateTracker( ExecutionStateSampler.newForTest(), diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java index 679227a11dc0..27b11ad67c6d 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourcesTest.java @@ -103,6 +103,7 @@ import org.apache.beam.runners.dataflow.worker.windmill.Windmill; import org.apache.beam.runners.dataflow.worker.windmill.client.getdata.FakeGetDataClient; import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateCache; +import org.apache.beam.runners.dataflow.worker.windmill.state.WindmillStateReader; import org.apache.beam.runners.dataflow.worker.windmill.work.processing.failures.FailureTracker; import org.apache.beam.runners.dataflow.worker.windmill.work.refresh.HeartbeatSender; import org.apache.beam.sdk.Pipeline; @@ -1003,7 +1004,7 @@ public void testFailedWorkItemsAbort() throws Exception { WindmillStateCache.builder() .setSizeMb(options.getWorkerCacheMb()) .build() - .forComputation(COMPUTATION_ID), + .forComputation(COMPUTATION_ID, "systemName"), StreamingStepMetricsContainer.createRegistry(), new DataflowExecutionStateTracker( ExecutionStateSampler.newForTest(), diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java index f57e20d4b5fb..6785ce47d0f6 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCacheTest.java @@ -86,7 +86,10 @@ private static ExecutableWork createWork(ShardedKey shardedKey, long workToken, public void setUp() { computationStateCache = ComputationStateCache.create( - configFetcher, workExecutor, ignored -> stateCache, IdGenerators.decrementingLongs()); + configFetcher, + workExecutor, + (ignored1, ignored2) -> stateCache, + IdGenerators.decrementingLongs()); } @Test diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java index 87b746089f11..0b55a5119564 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternalsTest.java @@ -225,7 +225,7 @@ public void resetUnderTest() { mockReader, false, cache - .forComputation("comp") + .forComputation("comp", "systemName") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyKey", StandardCharsets.UTF_8), 123), @@ -241,7 +241,7 @@ public void resetUnderTest() { mockReader, true, cache - .forComputation("comp") + .forComputation("comp", "systemName") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyNewKey", StandardCharsets.UTF_8), 123), @@ -257,7 +257,7 @@ public void resetUnderTest() { mockReader, false, cacheViaMultimap - .forComputation("comp") + .forComputation("comp", "systemName") .forKey( WindmillComputationKey.create( "comp", ByteString.copyFrom("dummyNewKey", StandardCharsets.UTF_8), 123), @@ -2049,7 +2049,7 @@ false, key(NAMESPACE, tag), STATE_FAMILY, VarIntCoder.of())) // clear cache and recreate multimapState cache - .forComputation("comp") + .forComputation("comp", "systemName") .invalidate(ByteString.copyFrom("dummyKey", StandardCharsets.UTF_8), 123); resetUnderTest(); multimapState = underTest.state(NAMESPACE, addr); diff --git a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java index caa25bf83090..e711a780a4dc 100644 --- a/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java +++ b/runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresherTest.java @@ -257,7 +257,7 @@ public void testInvalidateStuckCommits() throws InterruptedException { ByteString key = ByteString.EMPTY; for (int i = 0; i < 5; i++) { WindmillStateCache.ForComputation perComputationStateCache = - spy(stateCache.forComputation(COMPUTATION_ID_PREFIX + i)); + spy(stateCache.forComputation(COMPUTATION_ID_PREFIX + i, "systemName" + i)); ComputationState computationState = spy(createComputationState(i, perComputationStateCache)); ExecutableWork fakeWork = createOldWork(ShardedKey.create(key, i), i, ignored -> {}); fakeWork.work().setState(Work.State.COMMITTING); diff --git a/sdks/go.mod b/sdks/go.mod index 9a8f5495acbd..c2e62141d6b8 100644 --- a/sdks/go.mod +++ b/sdks/go.mod @@ -26,8 +26,8 @@ toolchain go1.26.2 require ( cloud.google.com/go/bigquery v1.79.0 - cloud.google.com/go/bigtable v1.51.0 - cloud.google.com/go/datastore v1.26.0 + cloud.google.com/go/bigtable v1.50.0 + cloud.google.com/go/datastore v1.25.0 cloud.google.com/go/profiler v0.6.0 cloud.google.com/go/pubsub v1.51.0 cloud.google.com/go/spanner v1.94.0 @@ -35,8 +35,8 @@ require ( github.com/aws/aws-sdk-go-v2 v1.43.2 github.com/aws/aws-sdk-go-v2/config v1.32.33 github.com/aws/aws-sdk-go-v2/credentials v1.19.32 - github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.37 - github.com/aws/aws-sdk-go-v2/service/s3 v1.106.2 + github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.36 + github.com/aws/aws-sdk-go-v2/service/s3 v1.106.1 github.com/aws/smithy-go v1.27.5 github.com/docker/go-connections v0.7.0 // indirect github.com/dustin/go-humanize v1.0.1 @@ -153,9 +153,15 @@ require ( github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.33 // indirect github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.34 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.14 // indirect +<<<<<<< HEAD github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.26 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.33 // indirect github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.34 // indirect +======= + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.25 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.33 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.33 // indirect +>>>>>>> c6e0f5b630e (Bump github.com/aws/aws-sdk-go-v2/config in /sdks (#39551)) github.com/aws/aws-sdk-go-v2/service/sso v1.33.2 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.38.2 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.45.2 // indirect diff --git a/sdks/go.sum b/sdks/go.sum index 292e3948013b..854ef0472842 100644 --- a/sdks/go.sum +++ b/sdks/go.sum @@ -58,8 +58,8 @@ cloud.google.com/go/datacatalog v1.32.0 h1:fyYn8ODkGil5y3zTIqgIhOfzTu1ACaU2o+C75 cloud.google.com/go/datacatalog v1.32.0/go.mod h1:DE272tynQUwheJeQAyVfV+nO8yrdkuDyOgH2LtOrkWM= cloud.google.com/go/datastore v1.0.0/go.mod h1:LXYbyblFSglQ5pkeyhO+Qmw7ukd3C+pD7TKLgZqpHYE= cloud.google.com/go/datastore v1.1.0/go.mod h1:umbIZjpQpHh4hmRpGhH4tLFup+FVzqBi1b3c64qFpCk= -cloud.google.com/go/datastore v1.26.0 h1:9lgjj+DRv5Ay/tQ+vk9Ryz/G84ncnfwRC0RuHUGZm0U= -cloud.google.com/go/datastore v1.26.0/go.mod h1:jvJVNe+S2nHVIndV1H/B4s9K3MLsTMqOKlxSrzHTxB4= +cloud.google.com/go/datastore v1.25.0 h1:zUjMnCLCcRZVDSdQIXsbnNCl1SVRNw5Jm0J77gPaPKs= +cloud.google.com/go/datastore v1.25.0/go.mod h1:jvJVNe+S2nHVIndV1H/B4s9K3MLsTMqOKlxSrzHTxB4= cloud.google.com/go/firestore v1.6.1/go.mod h1:asNXNOzBdyVQmEU+ggO8UPodTkEVFW5Qx+rwHnAz+EY= cloud.google.com/go/iam v0.1.0/go.mod h1:vcUNEa0pEm0qRVpmWepWaFMIAI8/hjB9mO8rNCJtF6c= cloud.google.com/go/iam v0.1.1/go.mod h1:CKqrcnI/suGpybEHxZ7BMehL0oA4LpdyJdUlTl9jVMw= @@ -216,8 +216,8 @@ github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.33 h1:MobhiR6KIerWxmO74Zit5I github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.33/go.mod h1:xu02847OdZfNr/jAfZpHtyRk0b3v4d0kaoxNHxZGG/w= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.11.3/go.mod h1:0dHuD2HZZSiwfJSy1FO5bX1hQ1TxVV1QXXjpn3XUE44= github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.14.0/go.mod h1:UcgIwJ9KHquYxs6Q5skC9qXjhYMK+JASDYcXQ4X7JZE= -github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.37 h1:yOm5rq5yr2d5Kqu0GuRs4cThk8BW6ElvEvSfQ/bwOjk= -github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.37/go.mod h1:0hQ4udHxw6ioHPvG3euWtIqdW+NUk/gbGUYW8CfbijU= +github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.33 h1:T0FhDHSzJf4hcxzQv24E2Ul6dyFA3wQKmy8qFmzq85c= +github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.33/go.mod h1:SG4Q9PWeeNiaI5/SZt2OEQWtYJaqp48Gx9Gy9Fpkk9w= github.com/aws/aws-sdk-go-v2/internal/configsources v1.1.9/go.mod h1:AnVH5pvai0pAF4lXRq0bmhbes1u9R8wTE+g+183bZNM= github.com/aws/aws-sdk-go-v2/internal/configsources v1.2.3/go.mod h1:7sGSz1JCKHWWBHq98m6sMtWQikmYPpxjqOydDemiVoM= github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.33 h1:HAp1wLFZzch054uh3FK7rcVYg4v7J2FxVf3h3IGNZas= diff --git a/sdks/java/testing/nexmark/build.gradle b/sdks/java/testing/nexmark/build.gradle index 0eeaf931a889..768427d61bbf 100644 --- a/sdks/java/testing/nexmark/build.gradle +++ b/sdks/java/testing/nexmark/build.gradle @@ -128,6 +128,19 @@ def sparkJvmArgs() { return [] } +def sparkJvmArgs() { + def testJavaVer = project.findProperty('testJavaVersion') ? (project.property('testJavaVersion') as int) : JavaVersion.current().majorVersion.toInteger() + if (testJavaVer >= 17) { + return [ + "--add-opens=java.base/sun.nio.ch=ALL-UNNAMED", + "--add-opens=java.base/java.nio=ALL-UNNAMED", + "--add-opens=java.base/java.util=ALL-UNNAMED", + "--add-opens=java.base/java.lang.invoke=ALL-UNNAMED" + ] + } + return [] +} + def getNexmarkArgs = { def nexmarkArgsStr = project.findProperty(nexmarkArgsProperty) ?: "" def nexmarkArgsList = new ArrayList() diff --git a/sdks/java/testing/tpcds/build.gradle b/sdks/java/testing/tpcds/build.gradle index 60c2f8bfdd8b..ac76e459d5e7 100644 --- a/sdks/java/testing/tpcds/build.gradle +++ b/sdks/java/testing/tpcds/build.gradle @@ -118,6 +118,19 @@ def sparkJvmArgs() { return [] } +def sparkJvmArgs() { + def testJavaVer = project.findProperty('testJavaVersion') ? (project.property('testJavaVersion') as int) : JavaVersion.current().majorVersion.toInteger() + if (testJavaVer >= 17) { + return [ + "--add-opens=java.base/sun.nio.ch=ALL-UNNAMED", + "--add-opens=java.base/java.nio=ALL-UNNAMED", + "--add-opens=java.base/java.util=ALL-UNNAMED", + "--add-opens=java.base/java.lang.invoke=ALL-UNNAMED" + ] + } + return [] +} + // Execute the TPC-DS queries or suites via Gradle. // // Parameters: