Skip to content

ARROW-17762: [C++] WIP: Add ordering information to Acero - #14158

Closed
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering
Closed

ARROW-17762: [C++] WIP: Add ordering information to Acero#14158
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering

Conversation

@westonpace

Copy link
Copy Markdown
Member

No description provided.

@github-actions

Copy link
Copy Markdown

@github-actions

Copy link
Copy Markdown

⚠️ Ticket has not been started in JIRA, please click 'Start Progress'.

@westonpace
westonpaceforce-pushed the feature/ARROW-17762--add-ordering branch from caf4a9f to 467f26fCompareSeptember 19, 2022 12:03
@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw FYI, I think you were interested in my plans for ordered execution. This PR is based on my earlier proposal I sent to the ML. I plan to create an example fetch node that consumes the ordering information to do something useful today.

This PR is built on top of ARROW-17287 and so it is a little easier to look at just the diff between the two:

westonpace/arrow@feature/ARROW-17287--initial-exec-plan-scan-node...westonpace:arrow:feature/ARROW-17762--add-ordering

@rtpsw

Copy link
Copy Markdown
Contributor

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw I've added an example of a FetchNode that consumes the ordering information. This will hopefully give some idea on how to use the exec batch index.

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

There is no rush. I probably won't get back to this myself for a while as I need to get #13782 (and numerous follow-ups) merged in.

@icexelloss

Copy link
Copy Markdown
Contributor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

@westonpace

Copy link
Copy Markdown
MemberAuthor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

Yes. If a node relies on ordering then it will resequence the batches before processing them. I try and take care to use both "reorder" and "resequence" independently as there are two rather different problems.

The first problem is when the input has no known ordering or is in a completely random order. In that case we must "reorder" which is "not streaming" and a "pipeline breaker" and requires us to cache all data in memory (or spill) in order to assign the order.

The second problem is when the input is mostly ordered but might be a bit noisy due to something like a parallel scan. In that case we already have a sequence number and we assume the sequence number is, generally, within some max delta from the correct ordering. In that case we only need to resequence (not reorder). This operation is "mostly streaming" and only sometimes a "pipeline breaker".

@zifengyu

Copy link
Copy Markdown

This feature is exactly what we need to adapt Acero. I tried to add ExecBatch ordering and implemented the limit operator in our product. Here is what we saw in the tests.

  1. It seems a little difficult to finish the node (and notify downstream node) as the input / output batch counts are not the same. In our case, the finish may happen either when having the limit number of rows or upstream node is finished producing (but not generated limit rows). The former occurs in Queue's deliver task while latter occurs in FetchNode's InputFinished. We did not find an easy way to sync these two components so we moved the queue part inside node and added a counter to track sent rows.

  2. We also need the offset setting to skip the first a few rows in the limit operator. Can this be included in FetchNode so we may switch back to Acero node in future?

Anyway, this proposal is critical to our using Acero. We are looking forward to its release.

@jorisvandenbossche

Copy link
Copy Markdown
Member

This is being closed in favor of other PRs as listed in the issue #32991

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@westonpace@rtpsw@icexelloss@zifengyu@jorisvandenbossche
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
ARROW-17762: [C++] WIP: Add ordering information to Acero by westonpace · Pull Request #14158 · apache/arrow · GitHub
Skip to content

ARROW-17762: [C++] WIP: Add ordering information to Acero - #14158

Closed
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering
Closed

ARROW-17762: [C++] WIP: Add ordering information to Acero#14158
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering

Conversation

@westonpace

Copy link
Copy Markdown
Member

No description provided.

@github-actions

Copy link
Copy Markdown

@github-actions

Copy link
Copy Markdown

⚠️ Ticket has not been started in JIRA, please click 'Start Progress'.

@westonpace
westonpaceforce-pushed the feature/ARROW-17762--add-ordering branch from caf4a9f to 467f26fCompareSeptember 19, 2022 12:03
@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw FYI, I think you were interested in my plans for ordered execution. This PR is based on my earlier proposal I sent to the ML. I plan to create an example fetch node that consumes the ordering information to do something useful today.

This PR is built on top of ARROW-17287 and so it is a little easier to look at just the diff between the two:

westonpace/arrow@feature/ARROW-17287--initial-exec-plan-scan-node...westonpace:arrow:feature/ARROW-17762--add-ordering

@rtpsw

Copy link
Copy Markdown
Contributor

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw I've added an example of a FetchNode that consumes the ordering information. This will hopefully give some idea on how to use the exec batch index.

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

There is no rush. I probably won't get back to this myself for a while as I need to get #13782 (and numerous follow-ups) merged in.

@icexelloss

Copy link
Copy Markdown
Contributor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

@westonpace

Copy link
Copy Markdown
MemberAuthor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

Yes. If a node relies on ordering then it will resequence the batches before processing them. I try and take care to use both "reorder" and "resequence" independently as there are two rather different problems.

The first problem is when the input has no known ordering or is in a completely random order. In that case we must "reorder" which is "not streaming" and a "pipeline breaker" and requires us to cache all data in memory (or spill) in order to assign the order.

The second problem is when the input is mostly ordered but might be a bit noisy due to something like a parallel scan. In that case we already have a sequence number and we assume the sequence number is, generally, within some max delta from the correct ordering. In that case we only need to resequence (not reorder). This operation is "mostly streaming" and only sometimes a "pipeline breaker".

@zifengyu

Copy link
Copy Markdown

This feature is exactly what we need to adapt Acero. I tried to add ExecBatch ordering and implemented the limit operator in our product. Here is what we saw in the tests.

  1. It seems a little difficult to finish the node (and notify downstream node) as the input / output batch counts are not the same. In our case, the finish may happen either when having the limit number of rows or upstream node is finished producing (but not generated limit rows). The former occurs in Queue's deliver task while latter occurs in FetchNode's InputFinished. We did not find an easy way to sync these two components so we moved the queue part inside node and added a counter to track sent rows.

  2. We also need the offset setting to skip the first a few rows in the limit operator. Can this be included in FetchNode so we may switch back to Acero node in future?

Anyway, this proposal is critical to our using Acero. We are looking forward to its release.

@jorisvandenbossche

Copy link
Copy Markdown
Member

This is being closed in favor of other PRs as listed in the issue #32991

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@westonpace@rtpsw@icexelloss@zifengyu@jorisvandenbossche
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' ARROW-17762: [C++] WIP: Add ordering information to Acero by westonpace · Pull Request #14158 · apache/arrow · GitHub
Skip to content

ARROW-17762: [C++] WIP: Add ordering information to Acero - #14158

Closed
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering
Closed

ARROW-17762: [C++] WIP: Add ordering information to Acero#14158
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering

Conversation

@westonpace

Copy link
Copy Markdown
Member

No description provided.

@github-actions

Copy link
Copy Markdown

@github-actions

Copy link
Copy Markdown

⚠️ Ticket has not been started in JIRA, please click 'Start Progress'.

@westonpace
westonpaceforce-pushed the feature/ARROW-17762--add-ordering branch from caf4a9f to 467f26fCompareSeptember 19, 2022 12:03
@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw FYI, I think you were interested in my plans for ordered execution. This PR is based on my earlier proposal I sent to the ML. I plan to create an example fetch node that consumes the ordering information to do something useful today.

This PR is built on top of ARROW-17287 and so it is a little easier to look at just the diff between the two:

westonpace/arrow@feature/ARROW-17287--initial-exec-plan-scan-node...westonpace:arrow:feature/ARROW-17762--add-ordering

@rtpsw

Copy link
Copy Markdown
Contributor

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw I've added an example of a FetchNode that consumes the ordering information. This will hopefully give some idea on how to use the exec batch index.

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

There is no rush. I probably won't get back to this myself for a while as I need to get #13782 (and numerous follow-ups) merged in.

@icexelloss

Copy link
Copy Markdown
Contributor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

@westonpace

Copy link
Copy Markdown
MemberAuthor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

Yes. If a node relies on ordering then it will resequence the batches before processing them. I try and take care to use both "reorder" and "resequence" independently as there are two rather different problems.

The first problem is when the input has no known ordering or is in a completely random order. In that case we must "reorder" which is "not streaming" and a "pipeline breaker" and requires us to cache all data in memory (or spill) in order to assign the order.

The second problem is when the input is mostly ordered but might be a bit noisy due to something like a parallel scan. In that case we already have a sequence number and we assume the sequence number is, generally, within some max delta from the correct ordering. In that case we only need to resequence (not reorder). This operation is "mostly streaming" and only sometimes a "pipeline breaker".

@zifengyu

Copy link
Copy Markdown

This feature is exactly what we need to adapt Acero. I tried to add ExecBatch ordering and implemented the limit operator in our product. Here is what we saw in the tests.

  1. It seems a little difficult to finish the node (and notify downstream node) as the input / output batch counts are not the same. In our case, the finish may happen either when having the limit number of rows or upstream node is finished producing (but not generated limit rows). The former occurs in Queue's deliver task while latter occurs in FetchNode's InputFinished. We did not find an easy way to sync these two components so we moved the queue part inside node and added a counter to track sent rows.

  2. We also need the offset setting to skip the first a few rows in the limit operator. Can this be included in FetchNode so we may switch back to Acero node in future?

Anyway, this proposal is critical to our using Acero. We are looking forward to its release.

@jorisvandenbossche

Copy link
Copy Markdown
Member

This is being closed in favor of other PRs as listed in the issue #32991

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@westonpace@rtpsw@icexelloss@zifengyu@jorisvandenbossche
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' ARROW-17762: [C++] WIP: Add ordering information to Acero by westonpace · Pull Request #14158 · apache/arrow · GitHub
Skip to content

ARROW-17762: [C++] WIP: Add ordering information to Acero - #14158

Closed
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering
Closed

ARROW-17762: [C++] WIP: Add ordering information to Acero#14158
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering

Conversation

@westonpace

Copy link
Copy Markdown
Member

No description provided.

@github-actions

Copy link
Copy Markdown

@github-actions

Copy link
Copy Markdown

⚠️ Ticket has not been started in JIRA, please click 'Start Progress'.

@westonpace
westonpaceforce-pushed the feature/ARROW-17762--add-ordering branch from caf4a9f to 467f26fCompareSeptember 19, 2022 12:03
@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw FYI, I think you were interested in my plans for ordered execution. This PR is based on my earlier proposal I sent to the ML. I plan to create an example fetch node that consumes the ordering information to do something useful today.

This PR is built on top of ARROW-17287 and so it is a little easier to look at just the diff between the two:

westonpace/arrow@feature/ARROW-17287--initial-exec-plan-scan-node...westonpace:arrow:feature/ARROW-17762--add-ordering

@rtpsw

Copy link
Copy Markdown
Contributor

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw I've added an example of a FetchNode that consumes the ordering information. This will hopefully give some idea on how to use the exec batch index.

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

There is no rush. I probably won't get back to this myself for a while as I need to get #13782 (and numerous follow-ups) merged in.

@icexelloss

Copy link
Copy Markdown
Contributor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

@westonpace

Copy link
Copy Markdown
MemberAuthor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

Yes. If a node relies on ordering then it will resequence the batches before processing them. I try and take care to use both "reorder" and "resequence" independently as there are two rather different problems.

The first problem is when the input has no known ordering or is in a completely random order. In that case we must "reorder" which is "not streaming" and a "pipeline breaker" and requires us to cache all data in memory (or spill) in order to assign the order.

The second problem is when the input is mostly ordered but might be a bit noisy due to something like a parallel scan. In that case we already have a sequence number and we assume the sequence number is, generally, within some max delta from the correct ordering. In that case we only need to resequence (not reorder). This operation is "mostly streaming" and only sometimes a "pipeline breaker".

@zifengyu

Copy link
Copy Markdown

This feature is exactly what we need to adapt Acero. I tried to add ExecBatch ordering and implemented the limit operator in our product. Here is what we saw in the tests.

  1. It seems a little difficult to finish the node (and notify downstream node) as the input / output batch counts are not the same. In our case, the finish may happen either when having the limit number of rows or upstream node is finished producing (but not generated limit rows). The former occurs in Queue's deliver task while latter occurs in FetchNode's InputFinished. We did not find an easy way to sync these two components so we moved the queue part inside node and added a counter to track sent rows.

  2. We also need the offset setting to skip the first a few rows in the limit operator. Can this be included in FetchNode so we may switch back to Acero node in future?

Anyway, this proposal is critical to our using Acero. We are looking forward to its release.

@jorisvandenbossche

Copy link
Copy Markdown
Member

This is being closed in favor of other PRs as listed in the issue #32991

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@westonpace@rtpsw@icexelloss@zifengyu@jorisvandenbossche
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' ARROW-17762: [C++] WIP: Add ordering information to Acero by westonpace · Pull Request #14158 · apache/arrow · GitHub
Skip to content

ARROW-17762: [C++] WIP: Add ordering information to Acero - #14158

Closed
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering
Closed

ARROW-17762: [C++] WIP: Add ordering information to Acero#14158
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering

Conversation

@westonpace

Copy link
Copy Markdown
Member

No description provided.

@github-actions

Copy link
Copy Markdown

@github-actions

Copy link
Copy Markdown

⚠️ Ticket has not been started in JIRA, please click 'Start Progress'.

@westonpace
westonpaceforce-pushed the feature/ARROW-17762--add-ordering branch from caf4a9f to 467f26fCompareSeptember 19, 2022 12:03
@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw FYI, I think you were interested in my plans for ordered execution. This PR is based on my earlier proposal I sent to the ML. I plan to create an example fetch node that consumes the ordering information to do something useful today.

This PR is built on top of ARROW-17287 and so it is a little easier to look at just the diff between the two:

westonpace/arrow@feature/ARROW-17287--initial-exec-plan-scan-node...westonpace:arrow:feature/ARROW-17762--add-ordering

@rtpsw

Copy link
Copy Markdown
Contributor

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw I've added an example of a FetchNode that consumes the ordering information. This will hopefully give some idea on how to use the exec batch index.

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

There is no rush. I probably won't get back to this myself for a while as I need to get #13782 (and numerous follow-ups) merged in.

@icexelloss

Copy link
Copy Markdown
Contributor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

@westonpace

Copy link
Copy Markdown
MemberAuthor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

Yes. If a node relies on ordering then it will resequence the batches before processing them. I try and take care to use both "reorder" and "resequence" independently as there are two rather different problems.

The first problem is when the input has no known ordering or is in a completely random order. In that case we must "reorder" which is "not streaming" and a "pipeline breaker" and requires us to cache all data in memory (or spill) in order to assign the order.

The second problem is when the input is mostly ordered but might be a bit noisy due to something like a parallel scan. In that case we already have a sequence number and we assume the sequence number is, generally, within some max delta from the correct ordering. In that case we only need to resequence (not reorder). This operation is "mostly streaming" and only sometimes a "pipeline breaker".

@zifengyu

Copy link
Copy Markdown

This feature is exactly what we need to adapt Acero. I tried to add ExecBatch ordering and implemented the limit operator in our product. Here is what we saw in the tests.

  1. It seems a little difficult to finish the node (and notify downstream node) as the input / output batch counts are not the same. In our case, the finish may happen either when having the limit number of rows or upstream node is finished producing (but not generated limit rows). The former occurs in Queue's deliver task while latter occurs in FetchNode's InputFinished. We did not find an easy way to sync these two components so we moved the queue part inside node and added a counter to track sent rows.

  2. We also need the offset setting to skip the first a few rows in the limit operator. Can this be included in FetchNode so we may switch back to Acero node in future?

Anyway, this proposal is critical to our using Acero. We are looking forward to its release.

@jorisvandenbossche

Copy link
Copy Markdown
Member

This is being closed in favor of other PRs as listed in the issue #32991

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@westonpace@rtpsw@icexelloss@zifengyu@jorisvandenbossche
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' ARROW-17762: [C++] WIP: Add ordering information to Acero by westonpace · Pull Request #14158 · apache/arrow · GitHub
Skip to content

ARROW-17762: [C++] WIP: Add ordering information to Acero - #14158

Closed
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering
Closed

ARROW-17762: [C++] WIP: Add ordering information to Acero#14158
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering

Conversation

@westonpace

Copy link
Copy Markdown
Member

No description provided.

@github-actions

Copy link
Copy Markdown

@github-actions

Copy link
Copy Markdown

⚠️ Ticket has not been started in JIRA, please click 'Start Progress'.

@westonpace
westonpaceforce-pushed the feature/ARROW-17762--add-ordering branch from caf4a9f to 467f26fCompareSeptember 19, 2022 12:03
@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw FYI, I think you were interested in my plans for ordered execution. This PR is based on my earlier proposal I sent to the ML. I plan to create an example fetch node that consumes the ordering information to do something useful today.

This PR is built on top of ARROW-17287 and so it is a little easier to look at just the diff between the two:

westonpace/arrow@feature/ARROW-17287--initial-exec-plan-scan-node...westonpace:arrow:feature/ARROW-17762--add-ordering

@rtpsw

Copy link
Copy Markdown
Contributor

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw I've added an example of a FetchNode that consumes the ordering information. This will hopefully give some idea on how to use the exec batch index.

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

There is no rush. I probably won't get back to this myself for a while as I need to get #13782 (and numerous follow-ups) merged in.

@icexelloss

Copy link
Copy Markdown
Contributor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

@westonpace

Copy link
Copy Markdown
MemberAuthor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

Yes. If a node relies on ordering then it will resequence the batches before processing them. I try and take care to use both "reorder" and "resequence" independently as there are two rather different problems.

The first problem is when the input has no known ordering or is in a completely random order. In that case we must "reorder" which is "not streaming" and a "pipeline breaker" and requires us to cache all data in memory (or spill) in order to assign the order.

The second problem is when the input is mostly ordered but might be a bit noisy due to something like a parallel scan. In that case we already have a sequence number and we assume the sequence number is, generally, within some max delta from the correct ordering. In that case we only need to resequence (not reorder). This operation is "mostly streaming" and only sometimes a "pipeline breaker".

@zifengyu

Copy link
Copy Markdown

This feature is exactly what we need to adapt Acero. I tried to add ExecBatch ordering and implemented the limit operator in our product. Here is what we saw in the tests.

  1. It seems a little difficult to finish the node (and notify downstream node) as the input / output batch counts are not the same. In our case, the finish may happen either when having the limit number of rows or upstream node is finished producing (but not generated limit rows). The former occurs in Queue's deliver task while latter occurs in FetchNode's InputFinished. We did not find an easy way to sync these two components so we moved the queue part inside node and added a counter to track sent rows.

  2. We also need the offset setting to skip the first a few rows in the limit operator. Can this be included in FetchNode so we may switch back to Acero node in future?

Anyway, this proposal is critical to our using Acero. We are looking forward to its release.

@jorisvandenbossche

Copy link
Copy Markdown
Member

This is being closed in favor of other PRs as listed in the issue #32991

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@westonpace@rtpsw@icexelloss@zifengyu@jorisvandenbossche
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); })(); ARROW-17762: [C++] WIP: Add ordering information to Acero by westonpace · Pull Request #14158 · apache/arrow · GitHub
Skip to content

ARROW-17762: [C++] WIP: Add ordering information to Acero - #14158

Closed
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering
Closed

ARROW-17762: [C++] WIP: Add ordering information to Acero#14158
westonpace wants to merge 11 commits into
apache:mainfrom
westonpace:feature/ARROW-17762--add-ordering

Conversation

@westonpace

Copy link
Copy Markdown
Member

No description provided.

@github-actions

Copy link
Copy Markdown

@github-actions

Copy link
Copy Markdown

⚠️ Ticket has not been started in JIRA, please click 'Start Progress'.

@westonpace
westonpaceforce-pushed the feature/ARROW-17762--add-ordering branch from caf4a9f to 467f26fCompareSeptember 19, 2022 12:03
@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw FYI, I think you were interested in my plans for ordered execution. This PR is based on my earlier proposal I sent to the ML. I plan to create an example fetch node that consumes the ordering information to do something useful today.

This PR is built on top of ARROW-17287 and so it is a little easier to look at just the diff between the two:

westonpace/arrow@feature/ARROW-17287--initial-exec-plan-scan-node...westonpace:arrow:feature/ARROW-17762--add-ordering

@rtpsw

Copy link
Copy Markdown
Contributor

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

@westonpace

Copy link
Copy Markdown
MemberAuthor

@rtpsw I've added an example of a FetchNode that consumes the ordering information. This will hopefully give some idea on how to use the exec batch index.

Thanks, @westonpace. I'm interested though will need a couple of days to get to this.

There is no rush. I probably won't get back to this myself for a while as I need to get #13782 (and numerous follow-ups) merged in.

@icexelloss

Copy link
Copy Markdown
Contributor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

@westonpace

Copy link
Copy Markdown
MemberAuthor

I care about this work very much as well and hope can understand this better. If I remember correctly the high level idea is that there are nodes that requires ordering (e.g., asof join) and if the input batches are out of order (indicated by batch index), the consumer node will cache/reorder out of order batches before processing them?

Yes. If a node relies on ordering then it will resequence the batches before processing them. I try and take care to use both "reorder" and "resequence" independently as there are two rather different problems.

The first problem is when the input has no known ordering or is in a completely random order. In that case we must "reorder" which is "not streaming" and a "pipeline breaker" and requires us to cache all data in memory (or spill) in order to assign the order.

The second problem is when the input is mostly ordered but might be a bit noisy due to something like a parallel scan. In that case we already have a sequence number and we assume the sequence number is, generally, within some max delta from the correct ordering. In that case we only need to resequence (not reorder). This operation is "mostly streaming" and only sometimes a "pipeline breaker".

@zifengyu

Copy link
Copy Markdown

This feature is exactly what we need to adapt Acero. I tried to add ExecBatch ordering and implemented the limit operator in our product. Here is what we saw in the tests.

  1. It seems a little difficult to finish the node (and notify downstream node) as the input / output batch counts are not the same. In our case, the finish may happen either when having the limit number of rows or upstream node is finished producing (but not generated limit rows). The former occurs in Queue's deliver task while latter occurs in FetchNode's InputFinished. We did not find an easy way to sync these two components so we moved the queue part inside node and added a counter to track sent rows.

  2. We also need the offset setting to skip the first a few rows in the limit operator. Can this be included in FetchNode so we may switch back to Acero node in future?

Anyway, this proposal is critical to our using Acero. We are looking forward to its release.

@jorisvandenbossche

Copy link
Copy Markdown
Member

This is being closed in favor of other PRs as listed in the issue #32991

Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@westonpace@rtpsw@icexelloss@zifengyu@jorisvandenbossche