A very stupid naive parallel join merge, uses Key() as value to merge two different input channels source1 -> \ -- join --> source2 -> Distributed: source1a -> |-> -- join --> source1b -> key-hash-route |-> -- join --> reduce -> source1n -> |-> -- joi
(ctx *plan.Context, l, r TaskRunner, p *plan.JoinMerge)
| 128 | // source2n -> |-> -- join --> |
| 129 | // |
| 130 | func NewJoinNaiveMerge(ctx *plan.Context, l, r TaskRunner, p *plan.JoinMerge) *JoinMerge { |
| 131 | |
| 132 | m := &JoinMerge{ |
| 133 | TaskBase: NewTaskBase(ctx), |
| 134 | colIndex: p.ColIndex, |
| 135 | } |
| 136 | |
| 137 | m.ltask = l |
| 138 | m.rtask = r |
| 139 | m.leftStmt = p.LeftFrom |
| 140 | m.rightStmt = p.RightFrom |
| 141 | |
| 142 | return m |
| 143 | } |
| 144 | |
| 145 | func (m *JoinMerge) Run() error { |
| 146 | defer m.Ctx.Recover() |