File tree
5 files changed
+105
-48
lines changed- python/pyspark/sql/streaming
- sql/core/src
- main/scala/org/apache/spark/sql/execution/python/streaming
- test/scala/org/apache/spark/sql/execution/python/streaming
5 files changed
+105
-48
lines changedLines changed: 15 additions & 31 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
14 | 14 |
| |
15 | 15 |
| |
16 | 16 |
| |
17 |
| - | |
| 17 | + | |
18 | 18 |
| |
19 | 19 |
| |
20 |
| - | |
| 20 | + | |
21 | 21 |
| |
22 | 22 |
| |
23 | 23 |
| |
24 |
| - | |
25 |
| - | |
26 |
| - | |
27 | 24 |
| |
28 | 25 |
| |
29 | 26 |
| |
| |||
38 | 35 |
| |
39 | 36 |
| |
40 | 37 |
| |
41 |
| - | |
| 38 | + | |
42 | 39 |
| |
43 |
| - | |
| 40 | + | |
44 | 41 |
| |
45 | 42 |
| |
46 | 43 |
| |
| |||
67 | 64 |
| |
68 | 65 |
| |
69 | 66 |
| |
70 |
| - | |
| 67 | + | |
71 | 68 |
| |
72 |
| - | |
| 69 | + | |
73 | 70 |
| |
74 | 71 |
| |
75 | 72 |
| |
| |||
85 | 82 |
| |
86 | 83 |
| |
87 | 84 |
| |
88 |
| - | |
89 |
| - | |
90 |
| - | |
91 |
| - | |
92 |
| - | |
93 |
| - | |
94 |
| - | |
95 |
| - | |
96 |
| - | |
97 |
| - | |
98 |
| - | |
99 |
| - | |
100 |
| - | |
101 |
| - | |
| 85 | + | |
102 | 86 |
| |
103 | 87 |
| |
104 | 88 |
| |
105 | 89 |
| |
106 | 90 |
| |
107 |
| - | |
| 91 | + | |
108 | 92 |
| |
109 |
| - | |
| 93 | + | |
110 | 94 |
| |
111 |
| - | |
112 |
| - | |
113 |
| - | |
114 |
| - | |
| 95 | + | |
| 96 | + | |
| 97 | + | |
| 98 | + | |
115 | 99 |
| |
116 | 100 |
| |
117 | 101 |
| |
| |||
143 | 127 |
| |
144 | 128 |
| |
145 | 129 |
| |
146 |
| - | |
| 130 | + | |
147 | 131 |
| |
148 | 132 |
| |
149 | 133 |
| |
| |||
160 | 144 |
| |
161 | 145 |
| |
162 | 146 |
| |
163 |
| - | |
| 147 | + | |
164 | 148 |
| |
165 | 149 |
| |
166 | 150 |
| |
|
Lines changed: 20 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
467 | 467 |
| |
468 | 468 |
| |
469 | 469 |
| |
| 470 | + | |
| 471 | + | |
| 472 | + | |
| 473 | + | |
| 474 | + | |
| 475 | + | |
| 476 | + | |
| 477 | + | |
| 478 | + | |
| 479 | + | |
| 480 | + | |
| 481 | + | |
| 482 | + | |
| 483 | + | |
| 484 | + | |
| 485 | + | |
| 486 | + | |
| 487 | + | |
| 488 | + | |
| 489 | + | |
470 | 490 |
| |
471 | 491 |
| |
472 | 492 |
| |
|
Lines changed: 21 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
26 | 26 |
| |
27 | 27 |
| |
28 | 28 |
| |
| 29 | + | |
29 | 30 |
| |
30 | 31 |
| |
31 | 32 |
| |
| |||
57 | 58 |
| |
58 | 59 |
| |
59 | 60 |
| |
| 61 | + | |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
| 76 | + | |
| 77 | + | |
| 78 | + | |
| 79 | + | |
| 80 | + | |
60 | 81 |
|
Lines changed: 25 additions & 5 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
475 | 475 |
| |
476 | 476 |
| |
477 | 477 |
| |
478 |
| - | |
| 478 | + | |
479 | 479 |
| |
480 | 480 |
| |
481 | 481 |
| |
| |||
487 | 487 |
| |
488 | 488 |
| |
489 | 489 |
| |
490 |
| - | |
491 | 490 |
| |
492 | 491 |
| |
| 492 | + | |
493 | 493 |
| |
494 |
| - | |
495 |
| - | |
496 | 494 |
| |
497 | 495 |
| |
498 | 496 |
| |
499 | 497 |
| |
500 | 498 |
| |
501 | 499 |
| |
502 | 500 |
| |
503 |
| - | |
| 501 | + | |
504 | 502 |
| |
505 | 503 |
| |
506 | 504 |
| |
| |||
511 | 509 |
| |
512 | 510 |
| |
513 | 511 |
| |
| 512 | + | |
| 513 | + | |
| 514 | + | |
| 515 | + | |
| 516 | + | |
| 517 | + | |
| 518 | + | |
| 519 | + | |
| 520 | + | |
| 521 | + | |
| 522 | + | |
| 523 | + | |
| 524 | + | |
| 525 | + | |
| 526 | + | |
| 527 | + | |
| 528 | + | |
| 529 | + | |
| 530 | + | |
| 531 | + | |
| 532 | + | |
| 533 | + | |
514 | 534 |
| |
515 | 535 |
| |
516 | 536 |
| |
|
Lines changed: 24 additions & 12 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
103 | 103 |
| |
104 | 104 |
| |
105 | 105 |
| |
| 106 | + | |
| 107 | + | |
106 | 108 |
| |
107 | 109 |
| |
108 | 110 |
| |
| |||
260 | 262 |
| |
261 | 263 |
| |
262 | 264 |
| |
263 |
| - | |
264 |
| - | |
| 265 | + | |
| 266 | + | |
| 267 | + | |
| 268 | + | |
265 | 269 |
| |
266 | 270 |
| |
267 | 271 |
| |
| |||
278 | 282 |
| |
279 | 283 |
| |
280 | 284 |
| |
281 |
| - | |
282 |
| - | |
| 285 | + | |
| 286 | + | |
| 287 | + | |
| 288 | + | |
283 | 289 |
| |
284 | 290 |
| |
285 | 291 |
| |
286 |
| - | |
287 |
| - | |
288 |
| - | |
289 |
| - | |
| 292 | + | |
| 293 | + | |
| 294 | + | |
| 295 | + | |
| 296 | + | |
| 297 | + | |
| 298 | + | |
290 | 299 |
| |
291 | 300 |
| |
292 | 301 |
| |
| |||
302 | 311 |
| |
303 | 312 |
| |
304 | 313 |
| |
| 314 | + | |
305 | 315 |
| |
306 | 316 |
| |
307 |
| - | |
308 |
| - | |
| 317 | + | |
| 318 | + | |
| 319 | + | |
| 320 | + | |
309 | 321 |
| |
310 | 322 |
| |
311 | 323 |
| |
312 | 324 |
| |
313 | 325 |
| |
314 | 326 |
| |
315 |
| - | |
| 327 | + | |
316 | 328 |
| |
317 | 329 |
| |
318 | 330 |
| |
| |||
328 | 340 |
| |
329 | 341 |
| |
330 | 342 |
| |
331 |
| - | |
| 343 | + | |
332 | 344 |
| |
333 | 345 |
| |
334 | 346 |
| |
|
0 commit comments