diff --git a/runners/prism/java/build.gradle b/runners/prism/java/build.gradle index 5e5ddbe139ee..c8898b48b718 100644 --- a/runners/prism/java/build.gradle +++ b/runners/prism/java/build.gradle @@ -106,14 +106,13 @@ def sickbayTests = [ // Prism doesn't support multiple TestStreams. 'org.apache.beam.sdk.testing.TestStreamTest.testMultipleStreams', - // Sometimes fails missing a final 'AFTER'. Otherwise, Hangs in ElementManager.FailBundle due to a held stageState lock. - 'org.apache.beam.sdk.testing.TestStreamTest.testMultiStage', // GroupIntoBatchesTest tests that fail: - // Teststream has bad KV encodings due to using an outer context. + // Wrong number of elements in windows after GroupIntoBatches: null + // Incorrect output collection after GroupIntoBatches: null 'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testInStreamingMode', 'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testBufferingTimerInFixedWindow', - // sdk worker disconnected + // varint overflow 4120879041988 on counting step 'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testBufferingTimerInGlobalWindow', // ShardedKey not yet implemented. 'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testWithShardedKeyInGlobalWindow', diff --git a/sdks/go/pkg/beam/runners/prism/internal/execute.go b/sdks/go/pkg/beam/runners/prism/internal/execute.go index cad1fb7e5479..e87d0507af28 100644 --- a/sdks/go/pkg/beam/runners/prism/internal/execute.go +++ b/sdks/go/pkg/beam/runners/prism/internal/execute.go @@ -285,8 +285,8 @@ func executePipeline(ctx context.Context, wks map[string]*worker.W, j *jobservic //slog.Warn("teststream bytes", "value", string(v), "bytes", v) return v } - // Hack for Java Strings in test stream, since it doesn't encode them correctly. - forceLP := cID == "StringUtf8Coder" || cID != pyld.GetCoderId() + // the coder from teststream payload has to be LP'ed + forceLP := cID != pyld.GetCoderId() if forceLP { // slog.Warn("recoding TestStreamValue", "cID", cID, "newUrn", coders[cID].GetSpec().GetUrn(), "payloadCoder", pyld.GetCoderId(), "oldUrn", coders[pyld.GetCoderId()].GetSpec().GetUrn()) // The coder needed length prefixing. For simplicity, add a length prefix to each diff --git a/sdks/go/pkg/beam/runners/prism/internal/preprocess.go b/sdks/go/pkg/beam/runners/prism/internal/preprocess.go index 3311bcced9f4..16fe6a095e68 100644 --- a/sdks/go/pkg/beam/runners/prism/internal/preprocess.go +++ b/sdks/go/pkg/beam/runners/prism/internal/preprocess.go @@ -27,6 +27,7 @@ import ( "github.com/apache/beam/sdks/v2/go/pkg/beam/runners/prism/internal/jobservices" "github.com/apache/beam/sdks/v2/go/pkg/beam/runners/prism/internal/urns" "golang.org/x/exp/maps" + "google.golang.org/protobuf/encoding/prototext" "google.golang.org/protobuf/proto" ) @@ -152,6 +153,7 @@ func (p *preprocessor) preProcessGraph(comps *pipepb.Components, j *jobservices. } // Extract URNs for the given transform. + slog.Debug("components in pipeline proto", "proto", prototext.Format(comps)) keptLeaves := maps.Keys(leaves) sort.Strings(keptLeaves) diff --git a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/TestStreamTranslation.java b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/TestStreamTranslation.java index db1a2f875c90..9de85778e3a4 100644 --- a/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/TestStreamTranslation.java +++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/TestStreamTranslation.java @@ -123,7 +123,8 @@ static RunnerApi.TestStreamPayload.Event eventToProto( .setTimestamp(element.getTimestamp().getMillis()) .setEncodedElement( ByteString.copyFrom( - CoderUtils.encodeToByteArray(coder, element.getValue())))); + CoderUtils.encodeToByteArray( + coder, element.getValue(), Coder.Context.NESTED)))); } return RunnerApi.TestStreamPayload.Event.newBuilder().setElementEvent(builder).build(); default: @@ -149,7 +150,8 @@ static TestStream.Event eventFromProto( protoEvent.getElementEvent().getElementsList()) { decodedElements.add( TimestampedValue.of( - CoderUtils.decodeFromByteArray(coder, element.getEncodedElement().toByteArray()), + CoderUtils.decodeFromByteArray( + coder, element.getEncodedElement().toByteArray(), Coder.Context.NESTED), new Instant(element.getTimestamp()))); } return TestStream.ElementEvent.add(decodedElements);