Skip to content

Commit

Permalink
Fix intersection creating many duplicate events (#359)
Browse files Browse the repository at this point in the history
  • Loading branch information
Ulimo authored Feb 11, 2024
1 parent c001f0f commit f446fad
Show file tree
Hide file tree
Showing 2 changed files with 47 additions and 1 deletion.
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,7 @@ private async IAsyncEnumerable<KeyValuePair<int, IStreamEvent>> OnLockingPrepare

public void CreateBlock()
{
singleReadSource = NoReadSourceInLoop();
singleReadSource = false;

_transformBlock = new TransformManyBlock<KeyValuePair<int, IStreamEvent>, KeyValuePair<int, IStreamEvent>>((r) =>
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,52 @@ public override Relation VisitZanzibarJoinIntersectWildcard(ZanzibarJoinIntersec
},
}
},
// l.object_type = r.object_type
new ScalarFunction()
{
ExtensionUri = FunctionsComparison.Uri,
ExtensionName = FunctionsComparison.Equal,
Arguments = new List<Expression>()
{
new DirectFieldReference()
{
ReferenceSegment = new StructReferenceSegment()
{
Field = ObjectTypeColumn
}
},
new DirectFieldReference()
{
ReferenceSegment = new StructReferenceSegment()
{
Field = left.OutputLength + ObjectTypeColumn
}
},
}
},
// l.object_id = r.object_id
new ScalarFunction()
{
ExtensionUri = FunctionsComparison.Uri,
ExtensionName = FunctionsComparison.Equal,
Arguments = new List<Expression>()
{
new DirectFieldReference()
{
ReferenceSegment = new StructReferenceSegment()
{
Field = ObjectIdColumn
}
},
new DirectFieldReference()
{
ReferenceSegment = new StructReferenceSegment()
{
Field = left.OutputLength + ObjectIdColumn
}
},
}
},
// l.user_id = '*'
new ScalarFunction()
{
Expand Down

0 comments on commit f446fad

Please sign in to comment.