Skip to content

Commit 8d2fe2b

Browse files
committed
Use nested context in java TestStream
1 parent 5485467 commit 8d2fe2b

4 files changed

Lines changed: 14 additions & 11 deletions

File tree

runners/prism/java/build.gradle

Lines changed: 6 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -106,15 +106,14 @@ def sickbayTests = [
106106

107107
// Prism doesn't support multiple TestStreams.
108108
'org.apache.beam.sdk.testing.TestStreamTest.testMultipleStreams',
109-
// Sometimes fails missing a final 'AFTER'. Otherwise, Hangs in ElementManager.FailBundle due to a held stageState lock.
110-
'org.apache.beam.sdk.testing.TestStreamTest.testMultiStage',
111109

112110
// GroupIntoBatchesTest tests that fail:
113-
// Teststream has bad KV encodings due to using an outer context.
114-
'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testInStreamingMode',
115-
'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testBufferingTimerInFixedWindow',
116-
// sdk worker disconnected
117-
'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testBufferingTimerInGlobalWindow',
111+
// Wrong number of elements in windows after GroupIntoBatches: null
112+
// Incorrect output collection after GroupIntoBatches: null
113+
// 'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testInStreamingMode',
114+
// 'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testBufferingTimerInFixedWindow',
115+
// varint overflow 4120879041988 on counting step
116+
// 'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testBufferingTimerInGlobalWindow',
118117
// ShardedKey not yet implemented.
119118
'org.apache.beam.sdk.transforms.GroupIntoBatchesTest.testWithShardedKeyInGlobalWindow',
120119

sdks/go/pkg/beam/runners/prism/internal/execute.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -285,8 +285,8 @@ func executePipeline(ctx context.Context, wks map[string]*worker.W, j *jobservic
285285
//slog.Warn("teststream bytes", "value", string(v), "bytes", v)
286286
return v
287287
}
288-
// Hack for Java Strings in test stream, since it doesn't encode them correctly.
289-
forceLP := cID == "StringUtf8Coder" || cID != pyld.GetCoderId()
288+
// the coder from teststream payload has to be LP'ed
289+
forceLP := cID != pyld.GetCoderId()
290290
if forceLP {
291291
// slog.Warn("recoding TestStreamValue", "cID", cID, "newUrn", coders[cID].GetSpec().GetUrn(), "payloadCoder", pyld.GetCoderId(), "oldUrn", coders[pyld.GetCoderId()].GetSpec().GetUrn())
292292
// The coder needed length prefixing. For simplicity, add a length prefix to each

sdks/go/pkg/beam/runners/prism/internal/preprocess.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@ import (
2727
"github.com/apache/beam/sdks/v2/go/pkg/beam/runners/prism/internal/jobservices"
2828
"github.com/apache/beam/sdks/v2/go/pkg/beam/runners/prism/internal/urns"
2929
"golang.org/x/exp/maps"
30+
"google.golang.org/protobuf/encoding/prototext"
3031
"google.golang.org/protobuf/proto"
3132
)
3233

@@ -152,6 +153,7 @@ func (p *preprocessor) preProcessGraph(comps *pipepb.Components, j *jobservices.
152153
}
153154

154155
// Extract URNs for the given transform.
156+
slog.Debug("components in pipeline proto", "proto", prototext.Format(comps))
155157

156158
keptLeaves := maps.Keys(leaves)
157159
sort.Strings(keptLeaves)

sdks/java/core/src/main/java/org/apache/beam/sdk/util/construction/TestStreamTranslation.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -123,7 +123,8 @@ static <T> RunnerApi.TestStreamPayload.Event eventToProto(
123123
.setTimestamp(element.getTimestamp().getMillis())
124124
.setEncodedElement(
125125
ByteString.copyFrom(
126-
CoderUtils.encodeToByteArray(coder, element.getValue()))));
126+
CoderUtils.encodeToByteArray(
127+
coder, element.getValue(), Coder.Context.NESTED))));
127128
}
128129
return RunnerApi.TestStreamPayload.Event.newBuilder().setElementEvent(builder).build();
129130
default:
@@ -149,7 +150,8 @@ static <T> TestStream.Event<T> eventFromProto(
149150
protoEvent.getElementEvent().getElementsList()) {
150151
decodedElements.add(
151152
TimestampedValue.of(
152-
CoderUtils.decodeFromByteArray(coder, element.getEncodedElement().toByteArray()),
153+
CoderUtils.decodeFromByteArray(
154+
coder, element.getEncodedElement().toByteArray(), Coder.Context.NESTED),
153155
new Instant(element.getTimestamp())));
154156
}
155157
return TestStream.ElementEvent.add(decodedElements);

0 commit comments

Comments
 (0)