Skip to content

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters - #13799

Merged
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299
Sep 1, 2022
Merged

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters#13799
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299

Conversation

@marsupialtail

Copy link
Copy Markdown
Contributor

This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.

This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS.

The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment.

To test this, set up an i3.2xlarge instance on AWS:

import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)

For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do

s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)

You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.

@github-actions

Copy link
Copy Markdown

Thanks for opening a pull request!

If this is not a minor PR. Could you open an issue for this pull request on JIRA? https://issues.apache.org/jira/browse/ARROW

Opening JIRAs ahead of time contributes to the Openness of the Apache Arrow project.

Then could you also rename pull request title in the following format?

ARROW-${JIRA_ID}: [${COMPONENT}] ${SUMMARY}

or

MINOR: [${COMPONENT}] ${SUMMARY}

See also:

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. I took a quick pass at review.

Comment threadcpp/src/arrow/dataset/scanner.cc
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
_populate_builder(builder, columns=columns, filter=filter,
batch_size=batch_size, use_threads=use_threads,
batch_size=batch_size, batch_readahead=batch_readahead,
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

I don't think we need to specify this kwarg if we're just going to specify the default.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a Cython quirk. You have to specify all the arguments.

Comment threadpython/pyarrow/_dataset.pyx
@marsupialtailmarsupialtail changed the title Arrow 17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersArrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 4, 2022
@pitrou

Copy link
Copy Markdown
Member

@westonpace@bkietz Why exactly does ScannerBuilder allow setting the same things that can be set in ScanOptions?

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@bkietzbkietz changed the title Arrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 9, 2022
@github-actions

Copy link
Copy Markdown

@bkietz

bkietz commented Aug 9, 2022

Copy link
Copy Markdown
Member

@pitrou@westonpace IIUC, ScannerBuilder is at this point mostly a wrapper around a scan options. Once upon a time it was needed to mediate the difference between single threaded and async scanners and to guard construction of a dataset wrapping a record batch reader, but this becomes less and less necessary as more datasets functionality becomes subsumed by the compute engine. (for example, I'd say there's no longer a motivation to support constructing datasets from record batch readers since the compute engine can use them as sources directly.) In short, I think what you're observing is ScannerBuilder on a gentle walk toward deprecation

@westonpace

Copy link
Copy Markdown
Member

Yes, scanner builder is on its way out, I hope, as part of #13782 (well, probably a follow-up). At the moment it still serves a slight purpose in that the projection option is a little hard to specify and it is something of a thorn when it comes to augmented fields.

I also agree with your other point. We spent considerable effort at one point making various things look like a dataset because datasets were the primary interface to the compute engine (e.g. filtering & projection). The record batch reader example is a good example. I'd even go so far as to say the InMemoryDataset is probably superfluous and a better option in the future would be a "table_source" node. The scanner should be reserved for the case where you have multiple sources of data, with the same (or devolved versions of the same) schema.

All that being said, I don't think readahead is going away. However, in the near future (again, #13782) I was pondering if we should reframe readahead as "roughly how many bytes of data should the scanner attempt to read ahead" instead of "batch readahead and fragment readahead".

@marsupialtail

marsupialtail commented Aug 12, 2022

Copy link
Copy Markdown
ContributorAuthor

I believe this is ready to be merged. @pitrou@westonpace

@westonpace
westonpace self-requested a review August 15, 2022 18:31
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

There is a potential problem with this. You can't increase the fragment readahead by too much, or else the first batch will be significantly delayed. Not sure how much a problem this is though.

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A few grammatical suggestions but otherwise I think this is a good addition. I think this may change to bytes_readahead / fragment_readahead before the release but it will be nice to have this in place already.

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
marsupialtailand others added 4 commits August 19, 2022 16:06
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

Don't think the failed checks have anything to do with me.

@pitrou

Copy link
Copy Markdown
Member

Don't think the failed checks have anything to do with me.

Indeed, they don't.

@pitrou

Copy link
Copy Markdown
Member

@marsupialtail Would you like to address @westonpace 's suggestions? Then I think we're good to go.

marsupialtailand others added 3 commits August 25, 2022 14:18
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

OK. I commited all the changes. @pitrou@westonpace

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update, just two suggestions below.

Comment threadpython/pyarrow/_dataset.pyx
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

done

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thank you @marsupialtail !

@pitrou
pitrou merged commit ec7e250 into apache:masterSep 1, 2022
@ursabot

Copy link
Copy Markdown

Benchmark runs are scheduled for baseline = 46f38dc and contender = ec7e250. ec7e250 is a master commit associated with this PR. Results will be available as each benchmark for each run completes.
Conbench compare runs links:
[Finished ⬇️0.0% ⬆️0.0%] ec2-t3-xlarge-us-east-2
[Failed] test-mac-arm
[Failed ⬇️0.27% ⬆️0.0%] ursa-i9-9960x
[Finished ⬇️0.75% ⬆️0.11%] ursa-thinkcentre-m75q
Buildkite builds:
[Finished] ec7e250c ec2-t3-xlarge-us-east-2
[Failed] ec7e250c test-mac-arm
[Failed] ec7e250c ursa-i9-9960x
[Finished] ec7e250c ursa-thinkcentre-m75q
[Finished] 46f38dca ec2-t3-xlarge-us-east-2
[Failed] 46f38dca test-mac-arm
[Failed] 46f38dca ursa-i9-9960x
[Finished] 46f38dca ursa-thinkcentre-m75q
Supported benchmarks:
ec2-t3-xlarge-us-east-2: Supported benchmark langs: Python, R. Runs only benchmarks with cloud = True
test-mac-arm: Supported benchmark langs: C++, Python, R
ursa-i9-9960x: Supported benchmark langs: Python, R, JavaScript
ursa-thinkcentre-m75q: Supported benchmark langs: C++, Java

zagto pushed a commit to zagto/arrow that referenced this pull request Oct 7, 2022
…and kDefaultFragmentReadahead parameters (apache#13799)
This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.
This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS. The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment. To test this, set up an i3.2xlarge instance on AWS: ```
import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)
```
For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do
```
s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)
```
You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.
Authored-by: Ziheng Wang <zihengw@stanford.edu>
Signed-off-by: Antoine Pitrou <antoine@python.org>
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

@marsupialtail@pitrou@bkietz@westonpace@ursabot
, '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-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters by marsupialtail · Pull Request #13799 · apache/arrow · GitHub
Skip to content

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters - #13799

Merged
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299
Sep 1, 2022
Merged

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters#13799
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299

Conversation

@marsupialtail

Copy link
Copy Markdown
Contributor

This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.

This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS.

The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment.

To test this, set up an i3.2xlarge instance on AWS:

import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)

For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do

s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)

You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.

@github-actions

Copy link
Copy Markdown

Thanks for opening a pull request!

If this is not a minor PR. Could you open an issue for this pull request on JIRA? https://issues.apache.org/jira/browse/ARROW

Opening JIRAs ahead of time contributes to the Openness of the Apache Arrow project.

Then could you also rename pull request title in the following format?

ARROW-${JIRA_ID}: [${COMPONENT}] ${SUMMARY}

or

MINOR: [${COMPONENT}] ${SUMMARY}

See also:

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. I took a quick pass at review.

Comment threadcpp/src/arrow/dataset/scanner.cc
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
_populate_builder(builder, columns=columns, filter=filter,
batch_size=batch_size, use_threads=use_threads,
batch_size=batch_size, batch_readahead=batch_readahead,
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

I don't think we need to specify this kwarg if we're just going to specify the default.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a Cython quirk. You have to specify all the arguments.

Comment threadpython/pyarrow/_dataset.pyx
@marsupialtailmarsupialtail changed the title Arrow 17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersArrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 4, 2022
@pitrou

Copy link
Copy Markdown
Member

@westonpace@bkietz Why exactly does ScannerBuilder allow setting the same things that can be set in ScanOptions?

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@bkietzbkietz changed the title Arrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 9, 2022
@github-actions

Copy link
Copy Markdown

@bkietz

bkietz commented Aug 9, 2022

Copy link
Copy Markdown
Member

@pitrou@westonpace IIUC, ScannerBuilder is at this point mostly a wrapper around a scan options. Once upon a time it was needed to mediate the difference between single threaded and async scanners and to guard construction of a dataset wrapping a record batch reader, but this becomes less and less necessary as more datasets functionality becomes subsumed by the compute engine. (for example, I'd say there's no longer a motivation to support constructing datasets from record batch readers since the compute engine can use them as sources directly.) In short, I think what you're observing is ScannerBuilder on a gentle walk toward deprecation

@westonpace

Copy link
Copy Markdown
Member

Yes, scanner builder is on its way out, I hope, as part of #13782 (well, probably a follow-up). At the moment it still serves a slight purpose in that the projection option is a little hard to specify and it is something of a thorn when it comes to augmented fields.

I also agree with your other point. We spent considerable effort at one point making various things look like a dataset because datasets were the primary interface to the compute engine (e.g. filtering & projection). The record batch reader example is a good example. I'd even go so far as to say the InMemoryDataset is probably superfluous and a better option in the future would be a "table_source" node. The scanner should be reserved for the case where you have multiple sources of data, with the same (or devolved versions of the same) schema.

All that being said, I don't think readahead is going away. However, in the near future (again, #13782) I was pondering if we should reframe readahead as "roughly how many bytes of data should the scanner attempt to read ahead" instead of "batch readahead and fragment readahead".

@marsupialtail

marsupialtail commented Aug 12, 2022

Copy link
Copy Markdown
ContributorAuthor

I believe this is ready to be merged. @pitrou@westonpace

@westonpace
westonpace self-requested a review August 15, 2022 18:31
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

There is a potential problem with this. You can't increase the fragment readahead by too much, or else the first batch will be significantly delayed. Not sure how much a problem this is though.

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A few grammatical suggestions but otherwise I think this is a good addition. I think this may change to bytes_readahead / fragment_readahead before the release but it will be nice to have this in place already.

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
marsupialtailand others added 4 commits August 19, 2022 16:06
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

Don't think the failed checks have anything to do with me.

@pitrou

Copy link
Copy Markdown
Member

Don't think the failed checks have anything to do with me.

Indeed, they don't.

@pitrou

Copy link
Copy Markdown
Member

@marsupialtail Would you like to address @westonpace 's suggestions? Then I think we're good to go.

marsupialtailand others added 3 commits August 25, 2022 14:18
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

OK. I commited all the changes. @pitrou@westonpace

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update, just two suggestions below.

Comment threadpython/pyarrow/_dataset.pyx
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

done

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thank you @marsupialtail !

@pitrou
pitrou merged commit ec7e250 into apache:masterSep 1, 2022
@ursabot

Copy link
Copy Markdown

Benchmark runs are scheduled for baseline = 46f38dc and contender = ec7e250. ec7e250 is a master commit associated with this PR. Results will be available as each benchmark for each run completes.
Conbench compare runs links:
[Finished ⬇️0.0% ⬆️0.0%] ec2-t3-xlarge-us-east-2
[Failed] test-mac-arm
[Failed ⬇️0.27% ⬆️0.0%] ursa-i9-9960x
[Finished ⬇️0.75% ⬆️0.11%] ursa-thinkcentre-m75q
Buildkite builds:
[Finished] ec7e250c ec2-t3-xlarge-us-east-2
[Failed] ec7e250c test-mac-arm
[Failed] ec7e250c ursa-i9-9960x
[Finished] ec7e250c ursa-thinkcentre-m75q
[Finished] 46f38dca ec2-t3-xlarge-us-east-2
[Failed] 46f38dca test-mac-arm
[Failed] 46f38dca ursa-i9-9960x
[Finished] 46f38dca ursa-thinkcentre-m75q
Supported benchmarks:
ec2-t3-xlarge-us-east-2: Supported benchmark langs: Python, R. Runs only benchmarks with cloud = True
test-mac-arm: Supported benchmark langs: C++, Python, R
ursa-i9-9960x: Supported benchmark langs: Python, R, JavaScript
ursa-thinkcentre-m75q: Supported benchmark langs: C++, Java

zagto pushed a commit to zagto/arrow that referenced this pull request Oct 7, 2022
…and kDefaultFragmentReadahead parameters (apache#13799)
This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.
This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS. The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment. To test this, set up an i3.2xlarge instance on AWS: ```
import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)
```
For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do
```
s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)
```
You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.
Authored-by: Ziheng Wang <zihengw@stanford.edu>
Signed-off-by: Antoine Pitrou <antoine@python.org>
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

@marsupialtail@pitrou@bkietz@westonpace@ursabot
, '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-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters by marsupialtail · Pull Request #13799 · apache/arrow · GitHub
Skip to content

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters - #13799

Merged
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299
Sep 1, 2022
Merged

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters#13799
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299

Conversation

@marsupialtail

Copy link
Copy Markdown
Contributor

This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.

This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS.

The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment.

To test this, set up an i3.2xlarge instance on AWS:

import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)

For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do

s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)

You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.

@github-actions

Copy link
Copy Markdown

Thanks for opening a pull request!

If this is not a minor PR. Could you open an issue for this pull request on JIRA? https://issues.apache.org/jira/browse/ARROW

Opening JIRAs ahead of time contributes to the Openness of the Apache Arrow project.

Then could you also rename pull request title in the following format?

ARROW-${JIRA_ID}: [${COMPONENT}] ${SUMMARY}

or

MINOR: [${COMPONENT}] ${SUMMARY}

See also:

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. I took a quick pass at review.

Comment threadcpp/src/arrow/dataset/scanner.cc
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
_populate_builder(builder, columns=columns, filter=filter,
batch_size=batch_size, use_threads=use_threads,
batch_size=batch_size, batch_readahead=batch_readahead,
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

I don't think we need to specify this kwarg if we're just going to specify the default.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a Cython quirk. You have to specify all the arguments.

Comment threadpython/pyarrow/_dataset.pyx
@marsupialtailmarsupialtail changed the title Arrow 17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersArrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 4, 2022
@pitrou

Copy link
Copy Markdown
Member

@westonpace@bkietz Why exactly does ScannerBuilder allow setting the same things that can be set in ScanOptions?

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@bkietzbkietz changed the title Arrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 9, 2022
@github-actions

Copy link
Copy Markdown

@bkietz

bkietz commented Aug 9, 2022

Copy link
Copy Markdown
Member

@pitrou@westonpace IIUC, ScannerBuilder is at this point mostly a wrapper around a scan options. Once upon a time it was needed to mediate the difference between single threaded and async scanners and to guard construction of a dataset wrapping a record batch reader, but this becomes less and less necessary as more datasets functionality becomes subsumed by the compute engine. (for example, I'd say there's no longer a motivation to support constructing datasets from record batch readers since the compute engine can use them as sources directly.) In short, I think what you're observing is ScannerBuilder on a gentle walk toward deprecation

@westonpace

Copy link
Copy Markdown
Member

Yes, scanner builder is on its way out, I hope, as part of #13782 (well, probably a follow-up). At the moment it still serves a slight purpose in that the projection option is a little hard to specify and it is something of a thorn when it comes to augmented fields.

I also agree with your other point. We spent considerable effort at one point making various things look like a dataset because datasets were the primary interface to the compute engine (e.g. filtering & projection). The record batch reader example is a good example. I'd even go so far as to say the InMemoryDataset is probably superfluous and a better option in the future would be a "table_source" node. The scanner should be reserved for the case where you have multiple sources of data, with the same (or devolved versions of the same) schema.

All that being said, I don't think readahead is going away. However, in the near future (again, #13782) I was pondering if we should reframe readahead as "roughly how many bytes of data should the scanner attempt to read ahead" instead of "batch readahead and fragment readahead".

@marsupialtail

marsupialtail commented Aug 12, 2022

Copy link
Copy Markdown
ContributorAuthor

I believe this is ready to be merged. @pitrou@westonpace

@westonpace
westonpace self-requested a review August 15, 2022 18:31
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

There is a potential problem with this. You can't increase the fragment readahead by too much, or else the first batch will be significantly delayed. Not sure how much a problem this is though.

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A few grammatical suggestions but otherwise I think this is a good addition. I think this may change to bytes_readahead / fragment_readahead before the release but it will be nice to have this in place already.

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
marsupialtailand others added 4 commits August 19, 2022 16:06
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

Don't think the failed checks have anything to do with me.

@pitrou

Copy link
Copy Markdown
Member

Don't think the failed checks have anything to do with me.

Indeed, they don't.

@pitrou

Copy link
Copy Markdown
Member

@marsupialtail Would you like to address @westonpace 's suggestions? Then I think we're good to go.

marsupialtailand others added 3 commits August 25, 2022 14:18
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

OK. I commited all the changes. @pitrou@westonpace

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update, just two suggestions below.

Comment threadpython/pyarrow/_dataset.pyx
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

done

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thank you @marsupialtail !

@pitrou
pitrou merged commit ec7e250 into apache:masterSep 1, 2022
@ursabot

Copy link
Copy Markdown

Benchmark runs are scheduled for baseline = 46f38dc and contender = ec7e250. ec7e250 is a master commit associated with this PR. Results will be available as each benchmark for each run completes.
Conbench compare runs links:
[Finished ⬇️0.0% ⬆️0.0%] ec2-t3-xlarge-us-east-2
[Failed] test-mac-arm
[Failed ⬇️0.27% ⬆️0.0%] ursa-i9-9960x
[Finished ⬇️0.75% ⬆️0.11%] ursa-thinkcentre-m75q
Buildkite builds:
[Finished] ec7e250c ec2-t3-xlarge-us-east-2
[Failed] ec7e250c test-mac-arm
[Failed] ec7e250c ursa-i9-9960x
[Finished] ec7e250c ursa-thinkcentre-m75q
[Finished] 46f38dca ec2-t3-xlarge-us-east-2
[Failed] 46f38dca test-mac-arm
[Failed] 46f38dca ursa-i9-9960x
[Finished] 46f38dca ursa-thinkcentre-m75q
Supported benchmarks:
ec2-t3-xlarge-us-east-2: Supported benchmark langs: Python, R. Runs only benchmarks with cloud = True
test-mac-arm: Supported benchmark langs: C++, Python, R
ursa-i9-9960x: Supported benchmark langs: Python, R, JavaScript
ursa-thinkcentre-m75q: Supported benchmark langs: C++, Java

zagto pushed a commit to zagto/arrow that referenced this pull request Oct 7, 2022
…and kDefaultFragmentReadahead parameters (apache#13799)
This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.
This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS. The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment. To test this, set up an i3.2xlarge instance on AWS: ```
import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)
```
For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do
```
s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)
```
You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.
Authored-by: Ziheng Wang <zihengw@stanford.edu>
Signed-off-by: Antoine Pitrou <antoine@python.org>
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

@marsupialtail@pitrou@bkietz@westonpace@ursabot
, '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-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters by marsupialtail · Pull Request #13799 · apache/arrow · GitHub
Skip to content

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters - #13799

Merged
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299
Sep 1, 2022
Merged

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters#13799
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299

Conversation

@marsupialtail

Copy link
Copy Markdown
Contributor

This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.

This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS.

The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment.

To test this, set up an i3.2xlarge instance on AWS:

import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)

For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do

s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)

You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.

@github-actions

Copy link
Copy Markdown

Thanks for opening a pull request!

If this is not a minor PR. Could you open an issue for this pull request on JIRA? https://issues.apache.org/jira/browse/ARROW

Opening JIRAs ahead of time contributes to the Openness of the Apache Arrow project.

Then could you also rename pull request title in the following format?

ARROW-${JIRA_ID}: [${COMPONENT}] ${SUMMARY}

or

MINOR: [${COMPONENT}] ${SUMMARY}

See also:

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. I took a quick pass at review.

Comment threadcpp/src/arrow/dataset/scanner.cc
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
_populate_builder(builder, columns=columns, filter=filter,
batch_size=batch_size, use_threads=use_threads,
batch_size=batch_size, batch_readahead=batch_readahead,
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

I don't think we need to specify this kwarg if we're just going to specify the default.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a Cython quirk. You have to specify all the arguments.

Comment threadpython/pyarrow/_dataset.pyx
@marsupialtailmarsupialtail changed the title Arrow 17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersArrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 4, 2022
@pitrou

Copy link
Copy Markdown
Member

@westonpace@bkietz Why exactly does ScannerBuilder allow setting the same things that can be set in ScanOptions?

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@bkietzbkietz changed the title Arrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 9, 2022
@github-actions

Copy link
Copy Markdown

@bkietz

bkietz commented Aug 9, 2022

Copy link
Copy Markdown
Member

@pitrou@westonpace IIUC, ScannerBuilder is at this point mostly a wrapper around a scan options. Once upon a time it was needed to mediate the difference between single threaded and async scanners and to guard construction of a dataset wrapping a record batch reader, but this becomes less and less necessary as more datasets functionality becomes subsumed by the compute engine. (for example, I'd say there's no longer a motivation to support constructing datasets from record batch readers since the compute engine can use them as sources directly.) In short, I think what you're observing is ScannerBuilder on a gentle walk toward deprecation

@westonpace

Copy link
Copy Markdown
Member

Yes, scanner builder is on its way out, I hope, as part of #13782 (well, probably a follow-up). At the moment it still serves a slight purpose in that the projection option is a little hard to specify and it is something of a thorn when it comes to augmented fields.

I also agree with your other point. We spent considerable effort at one point making various things look like a dataset because datasets were the primary interface to the compute engine (e.g. filtering & projection). The record batch reader example is a good example. I'd even go so far as to say the InMemoryDataset is probably superfluous and a better option in the future would be a "table_source" node. The scanner should be reserved for the case where you have multiple sources of data, with the same (or devolved versions of the same) schema.

All that being said, I don't think readahead is going away. However, in the near future (again, #13782) I was pondering if we should reframe readahead as "roughly how many bytes of data should the scanner attempt to read ahead" instead of "batch readahead and fragment readahead".

@marsupialtail

marsupialtail commented Aug 12, 2022

Copy link
Copy Markdown
ContributorAuthor

I believe this is ready to be merged. @pitrou@westonpace

@westonpace
westonpace self-requested a review August 15, 2022 18:31
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

There is a potential problem with this. You can't increase the fragment readahead by too much, or else the first batch will be significantly delayed. Not sure how much a problem this is though.

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A few grammatical suggestions but otherwise I think this is a good addition. I think this may change to bytes_readahead / fragment_readahead before the release but it will be nice to have this in place already.

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
marsupialtailand others added 4 commits August 19, 2022 16:06
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

Don't think the failed checks have anything to do with me.

@pitrou

Copy link
Copy Markdown
Member

Don't think the failed checks have anything to do with me.

Indeed, they don't.

@pitrou

Copy link
Copy Markdown
Member

@marsupialtail Would you like to address @westonpace 's suggestions? Then I think we're good to go.

marsupialtailand others added 3 commits August 25, 2022 14:18
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

OK. I commited all the changes. @pitrou@westonpace

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update, just two suggestions below.

Comment threadpython/pyarrow/_dataset.pyx
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

done

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thank you @marsupialtail !

@pitrou
pitrou merged commit ec7e250 into apache:masterSep 1, 2022
@ursabot

Copy link
Copy Markdown

Benchmark runs are scheduled for baseline = 46f38dc and contender = ec7e250. ec7e250 is a master commit associated with this PR. Results will be available as each benchmark for each run completes.
Conbench compare runs links:
[Finished ⬇️0.0% ⬆️0.0%] ec2-t3-xlarge-us-east-2
[Failed] test-mac-arm
[Failed ⬇️0.27% ⬆️0.0%] ursa-i9-9960x
[Finished ⬇️0.75% ⬆️0.11%] ursa-thinkcentre-m75q
Buildkite builds:
[Finished] ec7e250c ec2-t3-xlarge-us-east-2
[Failed] ec7e250c test-mac-arm
[Failed] ec7e250c ursa-i9-9960x
[Finished] ec7e250c ursa-thinkcentre-m75q
[Finished] 46f38dca ec2-t3-xlarge-us-east-2
[Failed] 46f38dca test-mac-arm
[Failed] 46f38dca ursa-i9-9960x
[Finished] 46f38dca ursa-thinkcentre-m75q
Supported benchmarks:
ec2-t3-xlarge-us-east-2: Supported benchmark langs: Python, R. Runs only benchmarks with cloud = True
test-mac-arm: Supported benchmark langs: C++, Python, R
ursa-i9-9960x: Supported benchmark langs: Python, R, JavaScript
ursa-thinkcentre-m75q: Supported benchmark langs: C++, Java

zagto pushed a commit to zagto/arrow that referenced this pull request Oct 7, 2022
…and kDefaultFragmentReadahead parameters (apache#13799)
This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.
This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS. The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment. To test this, set up an i3.2xlarge instance on AWS: ```
import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)
```
For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do
```
s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)
```
You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.
Authored-by: Ziheng Wang <zihengw@stanford.edu>
Signed-off-by: Antoine Pitrou <antoine@python.org>
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

@marsupialtail@pitrou@bkietz@westonpace@ursabot
, '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-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters by marsupialtail · Pull Request #13799 · apache/arrow · GitHub
Skip to content

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters - #13799

Merged
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299
Sep 1, 2022
Merged

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters#13799
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299

Conversation

@marsupialtail

Copy link
Copy Markdown
Contributor

This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.

This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS.

The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment.

To test this, set up an i3.2xlarge instance on AWS:

import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)

For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do

s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)

You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.

@github-actions

Copy link
Copy Markdown

Thanks for opening a pull request!

If this is not a minor PR. Could you open an issue for this pull request on JIRA? https://issues.apache.org/jira/browse/ARROW

Opening JIRAs ahead of time contributes to the Openness of the Apache Arrow project.

Then could you also rename pull request title in the following format?

ARROW-${JIRA_ID}: [${COMPONENT}] ${SUMMARY}

or

MINOR: [${COMPONENT}] ${SUMMARY}

See also:

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. I took a quick pass at review.

Comment threadcpp/src/arrow/dataset/scanner.cc
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
_populate_builder(builder, columns=columns, filter=filter,
batch_size=batch_size, use_threads=use_threads,
batch_size=batch_size, batch_readahead=batch_readahead,
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

I don't think we need to specify this kwarg if we're just going to specify the default.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a Cython quirk. You have to specify all the arguments.

Comment threadpython/pyarrow/_dataset.pyx
@marsupialtailmarsupialtail changed the title Arrow 17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersArrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 4, 2022
@pitrou

Copy link
Copy Markdown
Member

@westonpace@bkietz Why exactly does ScannerBuilder allow setting the same things that can be set in ScanOptions?

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@bkietzbkietz changed the title Arrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 9, 2022
@github-actions

Copy link
Copy Markdown

@bkietz

bkietz commented Aug 9, 2022

Copy link
Copy Markdown
Member

@pitrou@westonpace IIUC, ScannerBuilder is at this point mostly a wrapper around a scan options. Once upon a time it was needed to mediate the difference between single threaded and async scanners and to guard construction of a dataset wrapping a record batch reader, but this becomes less and less necessary as more datasets functionality becomes subsumed by the compute engine. (for example, I'd say there's no longer a motivation to support constructing datasets from record batch readers since the compute engine can use them as sources directly.) In short, I think what you're observing is ScannerBuilder on a gentle walk toward deprecation

@westonpace

Copy link
Copy Markdown
Member

Yes, scanner builder is on its way out, I hope, as part of #13782 (well, probably a follow-up). At the moment it still serves a slight purpose in that the projection option is a little hard to specify and it is something of a thorn when it comes to augmented fields.

I also agree with your other point. We spent considerable effort at one point making various things look like a dataset because datasets were the primary interface to the compute engine (e.g. filtering & projection). The record batch reader example is a good example. I'd even go so far as to say the InMemoryDataset is probably superfluous and a better option in the future would be a "table_source" node. The scanner should be reserved for the case where you have multiple sources of data, with the same (or devolved versions of the same) schema.

All that being said, I don't think readahead is going away. However, in the near future (again, #13782) I was pondering if we should reframe readahead as "roughly how many bytes of data should the scanner attempt to read ahead" instead of "batch readahead and fragment readahead".

@marsupialtail

marsupialtail commented Aug 12, 2022

Copy link
Copy Markdown
ContributorAuthor

I believe this is ready to be merged. @pitrou@westonpace

@westonpace
westonpace self-requested a review August 15, 2022 18:31
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

There is a potential problem with this. You can't increase the fragment readahead by too much, or else the first batch will be significantly delayed. Not sure how much a problem this is though.

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A few grammatical suggestions but otherwise I think this is a good addition. I think this may change to bytes_readahead / fragment_readahead before the release but it will be nice to have this in place already.

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
marsupialtailand others added 4 commits August 19, 2022 16:06
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

Don't think the failed checks have anything to do with me.

@pitrou

Copy link
Copy Markdown
Member

Don't think the failed checks have anything to do with me.

Indeed, they don't.

@pitrou

Copy link
Copy Markdown
Member

@marsupialtail Would you like to address @westonpace 's suggestions? Then I think we're good to go.

marsupialtailand others added 3 commits August 25, 2022 14:18
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

OK. I commited all the changes. @pitrou@westonpace

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update, just two suggestions below.

Comment threadpython/pyarrow/_dataset.pyx
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

done

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thank you @marsupialtail !

@pitrou
pitrou merged commit ec7e250 into apache:masterSep 1, 2022
@ursabot

Copy link
Copy Markdown

Benchmark runs are scheduled for baseline = 46f38dc and contender = ec7e250. ec7e250 is a master commit associated with this PR. Results will be available as each benchmark for each run completes.
Conbench compare runs links:
[Finished ⬇️0.0% ⬆️0.0%] ec2-t3-xlarge-us-east-2
[Failed] test-mac-arm
[Failed ⬇️0.27% ⬆️0.0%] ursa-i9-9960x
[Finished ⬇️0.75% ⬆️0.11%] ursa-thinkcentre-m75q
Buildkite builds:
[Finished] ec7e250c ec2-t3-xlarge-us-east-2
[Failed] ec7e250c test-mac-arm
[Failed] ec7e250c ursa-i9-9960x
[Finished] ec7e250c ursa-thinkcentre-m75q
[Finished] 46f38dca ec2-t3-xlarge-us-east-2
[Failed] 46f38dca test-mac-arm
[Failed] 46f38dca ursa-i9-9960x
[Finished] 46f38dca ursa-thinkcentre-m75q
Supported benchmarks:
ec2-t3-xlarge-us-east-2: Supported benchmark langs: Python, R. Runs only benchmarks with cloud = True
test-mac-arm: Supported benchmark langs: C++, Python, R
ursa-i9-9960x: Supported benchmark langs: Python, R, JavaScript
ursa-thinkcentre-m75q: Supported benchmark langs: C++, Java

zagto pushed a commit to zagto/arrow that referenced this pull request Oct 7, 2022
…and kDefaultFragmentReadahead parameters (apache#13799)
This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.
This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS. The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment. To test this, set up an i3.2xlarge instance on AWS: ```
import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)
```
For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do
```
s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)
```
You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.
Authored-by: Ziheng Wang <zihengw@stanford.edu>
Signed-off-by: Antoine Pitrou <antoine@python.org>
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

@marsupialtail@pitrou@bkietz@westonpace@ursabot
, '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-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters by marsupialtail · Pull Request #13799 · apache/arrow · GitHub
Skip to content

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters - #13799

Merged
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299
Sep 1, 2022
Merged

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters#13799
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299

Conversation

@marsupialtail

Copy link
Copy Markdown
Contributor

This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.

This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS.

The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment.

To test this, set up an i3.2xlarge instance on AWS:

import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)

For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do

s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)

You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.

@github-actions

Copy link
Copy Markdown

Thanks for opening a pull request!

If this is not a minor PR. Could you open an issue for this pull request on JIRA? https://issues.apache.org/jira/browse/ARROW

Opening JIRAs ahead of time contributes to the Openness of the Apache Arrow project.

Then could you also rename pull request title in the following format?

ARROW-${JIRA_ID}: [${COMPONENT}] ${SUMMARY}

or

MINOR: [${COMPONENT}] ${SUMMARY}

See also:

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. I took a quick pass at review.

Comment threadcpp/src/arrow/dataset/scanner.cc
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
_populate_builder(builder, columns=columns, filter=filter,
batch_size=batch_size, use_threads=use_threads,
batch_size=batch_size, batch_readahead=batch_readahead,
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

I don't think we need to specify this kwarg if we're just going to specify the default.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a Cython quirk. You have to specify all the arguments.

Comment threadpython/pyarrow/_dataset.pyx
@marsupialtailmarsupialtail changed the title Arrow 17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersArrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 4, 2022
@pitrou

Copy link
Copy Markdown
Member

@westonpace@bkietz Why exactly does ScannerBuilder allow setting the same things that can be set in ScanOptions?

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@bkietzbkietz changed the title Arrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 9, 2022
@github-actions

Copy link
Copy Markdown

@bkietz

bkietz commented Aug 9, 2022

Copy link
Copy Markdown
Member

@pitrou@westonpace IIUC, ScannerBuilder is at this point mostly a wrapper around a scan options. Once upon a time it was needed to mediate the difference between single threaded and async scanners and to guard construction of a dataset wrapping a record batch reader, but this becomes less and less necessary as more datasets functionality becomes subsumed by the compute engine. (for example, I'd say there's no longer a motivation to support constructing datasets from record batch readers since the compute engine can use them as sources directly.) In short, I think what you're observing is ScannerBuilder on a gentle walk toward deprecation

@westonpace

Copy link
Copy Markdown
Member

Yes, scanner builder is on its way out, I hope, as part of #13782 (well, probably a follow-up). At the moment it still serves a slight purpose in that the projection option is a little hard to specify and it is something of a thorn when it comes to augmented fields.

I also agree with your other point. We spent considerable effort at one point making various things look like a dataset because datasets were the primary interface to the compute engine (e.g. filtering & projection). The record batch reader example is a good example. I'd even go so far as to say the InMemoryDataset is probably superfluous and a better option in the future would be a "table_source" node. The scanner should be reserved for the case where you have multiple sources of data, with the same (or devolved versions of the same) schema.

All that being said, I don't think readahead is going away. However, in the near future (again, #13782) I was pondering if we should reframe readahead as "roughly how many bytes of data should the scanner attempt to read ahead" instead of "batch readahead and fragment readahead".

@marsupialtail

marsupialtail commented Aug 12, 2022

Copy link
Copy Markdown
ContributorAuthor

I believe this is ready to be merged. @pitrou@westonpace

@westonpace
westonpace self-requested a review August 15, 2022 18:31
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

There is a potential problem with this. You can't increase the fragment readahead by too much, or else the first batch will be significantly delayed. Not sure how much a problem this is though.

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A few grammatical suggestions but otherwise I think this is a good addition. I think this may change to bytes_readahead / fragment_readahead before the release but it will be nice to have this in place already.

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
marsupialtailand others added 4 commits August 19, 2022 16:06
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

Don't think the failed checks have anything to do with me.

@pitrou

Copy link
Copy Markdown
Member

Don't think the failed checks have anything to do with me.

Indeed, they don't.

@pitrou

Copy link
Copy Markdown
Member

@marsupialtail Would you like to address @westonpace 's suggestions? Then I think we're good to go.

marsupialtailand others added 3 commits August 25, 2022 14:18
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

OK. I commited all the changes. @pitrou@westonpace

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update, just two suggestions below.

Comment threadpython/pyarrow/_dataset.pyx
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

done

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thank you @marsupialtail !

@pitrou
pitrou merged commit ec7e250 into apache:masterSep 1, 2022
@ursabot

Copy link
Copy Markdown

Benchmark runs are scheduled for baseline = 46f38dc and contender = ec7e250. ec7e250 is a master commit associated with this PR. Results will be available as each benchmark for each run completes.
Conbench compare runs links:
[Finished ⬇️0.0% ⬆️0.0%] ec2-t3-xlarge-us-east-2
[Failed] test-mac-arm
[Failed ⬇️0.27% ⬆️0.0%] ursa-i9-9960x
[Finished ⬇️0.75% ⬆️0.11%] ursa-thinkcentre-m75q
Buildkite builds:
[Finished] ec7e250c ec2-t3-xlarge-us-east-2
[Failed] ec7e250c test-mac-arm
[Failed] ec7e250c ursa-i9-9960x
[Finished] ec7e250c ursa-thinkcentre-m75q
[Finished] 46f38dca ec2-t3-xlarge-us-east-2
[Failed] 46f38dca test-mac-arm
[Failed] 46f38dca ursa-i9-9960x
[Finished] 46f38dca ursa-thinkcentre-m75q
Supported benchmarks:
ec2-t3-xlarge-us-east-2: Supported benchmark langs: Python, R. Runs only benchmarks with cloud = True
test-mac-arm: Supported benchmark langs: C++, Python, R
ursa-i9-9960x: Supported benchmark langs: Python, R, JavaScript
ursa-thinkcentre-m75q: Supported benchmark langs: C++, Java

zagto pushed a commit to zagto/arrow that referenced this pull request Oct 7, 2022
…and kDefaultFragmentReadahead parameters (apache#13799)
This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.
This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS. The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment. To test this, set up an i3.2xlarge instance on AWS: ```
import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)
```
For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do
```
s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)
```
You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.
Authored-by: Ziheng Wang <zihengw@stanford.edu>
Signed-off-by: Antoine Pitrou <antoine@python.org>
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

@marsupialtail@pitrou@bkietz@westonpace@ursabot
, '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); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters by marsupialtail · Pull Request #13799 · apache/arrow · GitHub
Skip to content

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters - #13799

Merged
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299
Sep 1, 2022
Merged

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters#13799
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299

Conversation

@marsupialtail

Copy link
Copy Markdown
Contributor

This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.

This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS.

The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment.

To test this, set up an i3.2xlarge instance on AWS:

import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)

For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do

s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)

You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.

@github-actions

Copy link
Copy Markdown

Thanks for opening a pull request!

If this is not a minor PR. Could you open an issue for this pull request on JIRA? https://issues.apache.org/jira/browse/ARROW

Opening JIRAs ahead of time contributes to the Openness of the Apache Arrow project.

Then could you also rename pull request title in the following format?

ARROW-${JIRA_ID}: [${COMPONENT}] ${SUMMARY}

or

MINOR: [${COMPONENT}] ${SUMMARY}

See also:

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. I took a quick pass at review.

Comment threadcpp/src/arrow/dataset/scanner.cc
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
_populate_builder(builder, columns=columns, filter=filter,
batch_size=batch_size, use_threads=use_threads,
batch_size=batch_size, batch_readahead=batch_readahead,
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

I don't think we need to specify this kwarg if we're just going to specify the default.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a Cython quirk. You have to specify all the arguments.

Comment threadpython/pyarrow/_dataset.pyx
@marsupialtailmarsupialtail changed the title Arrow 17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersArrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 4, 2022
@pitrou

Copy link
Copy Markdown
Member

@westonpace@bkietz Why exactly does ScannerBuilder allow setting the same things that can be set in ScanOptions?

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@bkietzbkietz changed the title Arrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 9, 2022
@github-actions

Copy link
Copy Markdown

@bkietz

bkietz commented Aug 9, 2022

Copy link
Copy Markdown
Member

@pitrou@westonpace IIUC, ScannerBuilder is at this point mostly a wrapper around a scan options. Once upon a time it was needed to mediate the difference between single threaded and async scanners and to guard construction of a dataset wrapping a record batch reader, but this becomes less and less necessary as more datasets functionality becomes subsumed by the compute engine. (for example, I'd say there's no longer a motivation to support constructing datasets from record batch readers since the compute engine can use them as sources directly.) In short, I think what you're observing is ScannerBuilder on a gentle walk toward deprecation

@westonpace

Copy link
Copy Markdown
Member

Yes, scanner builder is on its way out, I hope, as part of #13782 (well, probably a follow-up). At the moment it still serves a slight purpose in that the projection option is a little hard to specify and it is something of a thorn when it comes to augmented fields.

I also agree with your other point. We spent considerable effort at one point making various things look like a dataset because datasets were the primary interface to the compute engine (e.g. filtering & projection). The record batch reader example is a good example. I'd even go so far as to say the InMemoryDataset is probably superfluous and a better option in the future would be a "table_source" node. The scanner should be reserved for the case where you have multiple sources of data, with the same (or devolved versions of the same) schema.

All that being said, I don't think readahead is going away. However, in the near future (again, #13782) I was pondering if we should reframe readahead as "roughly how many bytes of data should the scanner attempt to read ahead" instead of "batch readahead and fragment readahead".

@marsupialtail

marsupialtail commented Aug 12, 2022

Copy link
Copy Markdown
ContributorAuthor

I believe this is ready to be merged. @pitrou@westonpace

@westonpace
westonpace self-requested a review August 15, 2022 18:31
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

There is a potential problem with this. You can't increase the fragment readahead by too much, or else the first batch will be significantly delayed. Not sure how much a problem this is though.

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A few grammatical suggestions but otherwise I think this is a good addition. I think this may change to bytes_readahead / fragment_readahead before the release but it will be nice to have this in place already.

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
marsupialtailand others added 4 commits August 19, 2022 16:06
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

Don't think the failed checks have anything to do with me.

@pitrou

Copy link
Copy Markdown
Member

Don't think the failed checks have anything to do with me.

Indeed, they don't.

@pitrou

Copy link
Copy Markdown
Member

@marsupialtail Would you like to address @westonpace 's suggestions? Then I think we're good to go.

marsupialtailand others added 3 commits August 25, 2022 14:18
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

OK. I commited all the changes. @pitrou@westonpace

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update, just two suggestions below.

Comment threadpython/pyarrow/_dataset.pyx
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

done

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thank you @marsupialtail !

@pitrou
pitrou merged commit ec7e250 into apache:masterSep 1, 2022
@ursabot

Copy link
Copy Markdown

Benchmark runs are scheduled for baseline = 46f38dc and contender = ec7e250. ec7e250 is a master commit associated with this PR. Results will be available as each benchmark for each run completes.
Conbench compare runs links:
[Finished ⬇️0.0% ⬆️0.0%] ec2-t3-xlarge-us-east-2
[Failed] test-mac-arm
[Failed ⬇️0.27% ⬆️0.0%] ursa-i9-9960x
[Finished ⬇️0.75% ⬆️0.11%] ursa-thinkcentre-m75q
Buildkite builds:
[Finished] ec7e250c ec2-t3-xlarge-us-east-2
[Failed] ec7e250c test-mac-arm
[Failed] ec7e250c ursa-i9-9960x
[Finished] ec7e250c ursa-thinkcentre-m75q
[Finished] 46f38dca ec2-t3-xlarge-us-east-2
[Failed] 46f38dca test-mac-arm
[Failed] 46f38dca ursa-i9-9960x
[Finished] 46f38dca ursa-thinkcentre-m75q
Supported benchmarks:
ec2-t3-xlarge-us-east-2: Supported benchmark langs: Python, R. Runs only benchmarks with cloud = True
test-mac-arm: Supported benchmark langs: C++, Python, R
ursa-i9-9960x: Supported benchmark langs: Python, R, JavaScript
ursa-thinkcentre-m75q: Supported benchmark langs: C++, Java

zagto pushed a commit to zagto/arrow that referenced this pull request Oct 7, 2022
…and kDefaultFragmentReadahead parameters (apache#13799)
This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.
This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS. The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment. To test this, set up an i3.2xlarge instance on AWS: ```
import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)
```
For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do
```
s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)
```
You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.
Authored-by: Ziheng Wang <zihengw@stanford.edu>
Signed-off-by: Antoine Pitrou <antoine@python.org>
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

@marsupialtail@pitrou@bkietz@westonpace@ursabot
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters by marsupialtail · Pull Request #13799 · apache/arrow · GitHub
Skip to content

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters - #13799

Merged
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299
Sep 1, 2022
Merged

ARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parameters#13799
pitrou merged 16 commits into
apache:masterfrom
marsupialtail:JIRA-17299

Conversation

@marsupialtail

Copy link
Copy Markdown
Contributor

This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.

This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS.

The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment.

To test this, set up an i3.2xlarge instance on AWS:

import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)

For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do

s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)

You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.

@github-actions

Copy link
Copy Markdown

Thanks for opening a pull request!

If this is not a minor PR. Could you open an issue for this pull request on JIRA? https://issues.apache.org/jira/browse/ARROW

Opening JIRAs ahead of time contributes to the Openness of the Apache Arrow project.

Then could you also rename pull request title in the following format?

ARROW-${JIRA_ID}: [${COMPONENT}] ${SUMMARY}

or

MINOR: [${COMPONENT}] ${SUMMARY}

See also:

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for adding this. I took a quick pass at review.

Comment threadcpp/src/arrow/dataset/scanner.cc
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
_populate_builder(builder, columns=columns, filter=filter,
batch_size=batch_size, use_threads=use_threads,
batch_size=batch_size, batch_readahead=batch_readahead,
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
fragment_readahead=_DEFAULT_FRAGMENT_READAHEAD,

I don't think we need to specify this kwarg if we're just going to specify the default.

Copy link
Copy Markdown
ContributorAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This is a Cython quirk. You have to specify all the arguments.

Comment threadpython/pyarrow/_dataset.pyx
@marsupialtailmarsupialtail changed the title Arrow 17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersArrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 4, 2022
@pitrou

Copy link
Copy Markdown
Member

@westonpace@bkietz Why exactly does ScannerBuilder allow setting the same things that can be set in ScanOptions?

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@bkietzbkietz changed the title Arrow-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersARROW-17299: [C++][Python] Expose the Scanner kDefaultBatchReadahead and kDefaultFragmentReadahead parametersAug 9, 2022
@github-actions

Copy link
Copy Markdown

@bkietz

bkietz commented Aug 9, 2022

Copy link
Copy Markdown
Member

@pitrou@westonpace IIUC, ScannerBuilder is at this point mostly a wrapper around a scan options. Once upon a time it was needed to mediate the difference between single threaded and async scanners and to guard construction of a dataset wrapping a record batch reader, but this becomes less and less necessary as more datasets functionality becomes subsumed by the compute engine. (for example, I'd say there's no longer a motivation to support constructing datasets from record batch readers since the compute engine can use them as sources directly.) In short, I think what you're observing is ScannerBuilder on a gentle walk toward deprecation

@westonpace

Copy link
Copy Markdown
Member

Yes, scanner builder is on its way out, I hope, as part of #13782 (well, probably a follow-up). At the moment it still serves a slight purpose in that the projection option is a little hard to specify and it is something of a thorn when it comes to augmented fields.

I also agree with your other point. We spent considerable effort at one point making various things look like a dataset because datasets were the primary interface to the compute engine (e.g. filtering & projection). The record batch reader example is a good example. I'd even go so far as to say the InMemoryDataset is probably superfluous and a better option in the future would be a "table_source" node. The scanner should be reserved for the case where you have multiple sources of data, with the same (or devolved versions of the same) schema.

All that being said, I don't think readahead is going away. However, in the near future (again, #13782) I was pondering if we should reframe readahead as "roughly how many bytes of data should the scanner attempt to read ahead" instead of "batch readahead and fragment readahead".

@marsupialtail

marsupialtail commented Aug 12, 2022

Copy link
Copy Markdown
ContributorAuthor

I believe this is ready to be merged. @pitrou@westonpace

@westonpace
westonpace self-requested a review August 15, 2022 18:31
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

There is a potential problem with this. You can't increase the fragment readahead by too much, or else the first batch will be significantly delayed. Not sure how much a problem this is though.

@westonpacewestonpace left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A few grammatical suggestions but otherwise I think this is a good addition. I think this may change to bytes_readahead / fragment_readahead before the release but it will be nice to have this in place already.

Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
Comment threadpython/pyarrow/_dataset.pyx Outdated
marsupialtailand others added 4 commits August 19, 2022 16:06
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

Don't think the failed checks have anything to do with me.

@pitrou

Copy link
Copy Markdown
Member

Don't think the failed checks have anything to do with me.

Indeed, they don't.

@pitrou

Copy link
Copy Markdown
Member

@marsupialtail Would you like to address @westonpace 's suggestions? Then I think we're good to go.

marsupialtailand others added 3 commits August 25, 2022 14:18
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
Co-authored-by: Weston Pace <weston.pace@gmail.com>
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

OK. I commited all the changes. @pitrou@westonpace

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the update, just two suggestions below.

Comment threadpython/pyarrow/_dataset.pyx
Comment threadcpp/src/arrow/dataset/scanner.h Outdated
@marsupialtail

Copy link
Copy Markdown
ContributorAuthor

done

@pitroupitrou left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM. Thank you @marsupialtail !

@pitrou
pitrou merged commit ec7e250 into apache:masterSep 1, 2022
@ursabot

Copy link
Copy Markdown

Benchmark runs are scheduled for baseline = 46f38dc and contender = ec7e250. ec7e250 is a master commit associated with this PR. Results will be available as each benchmark for each run completes.
Conbench compare runs links:
[Finished ⬇️0.0% ⬆️0.0%] ec2-t3-xlarge-us-east-2
[Failed] test-mac-arm
[Failed ⬇️0.27% ⬆️0.0%] ursa-i9-9960x
[Finished ⬇️0.75% ⬆️0.11%] ursa-thinkcentre-m75q
Buildkite builds:
[Finished] ec7e250c ec2-t3-xlarge-us-east-2
[Failed] ec7e250c test-mac-arm
[Failed] ec7e250c ursa-i9-9960x
[Finished] ec7e250c ursa-thinkcentre-m75q
[Finished] 46f38dca ec2-t3-xlarge-us-east-2
[Failed] 46f38dca test-mac-arm
[Failed] 46f38dca ursa-i9-9960x
[Finished] 46f38dca ursa-thinkcentre-m75q
Supported benchmarks:
ec2-t3-xlarge-us-east-2: Supported benchmark langs: Python, R. Runs only benchmarks with cloud = True
test-mac-arm: Supported benchmark langs: C++, Python, R
ursa-i9-9960x: Supported benchmark langs: Python, R, JavaScript
ursa-thinkcentre-m75q: Supported benchmark langs: C++, Java

zagto pushed a commit to zagto/arrow that referenced this pull request Oct 7, 2022
…and kDefaultFragmentReadahead parameters (apache#13799)
This exposes the Fragment Readahead and Batch Readahead flags in the C++ Scanner to the user in Python.
This can be used to finetune RAM usage and IO utilization during downloading large files from S3 or other network sources. I believe the default settings are overly conservative for small RAM settings and I observe less than 20% IO utilization on some instances on AWS. The Python API is exposed only to methods where these flags make sense. Scanning from a RecordBatchIterator won't need those these flags nor will those flags make sense. Only the latter flag makes sense for making a scanner from a fragment. To test this, set up an i3.2xlarge instance on AWS: ```
import pyarrow
import pyarrow.dataset as ds
import pyarrow.csv as csv
import time
pyarrow.set_cpu_count(8)
pyarrow.set_io_thread_count(16)
lineitem_scheme = ["l_orderkey","l_partkey","l_suppkey","l_linenumber","l_quantity","l_extendedprice",
"l_discount","l_tax","l_returnflag","l_linestatus","l_shipdate","l_commitdate","l_receiptdate","l_shipinstruct",
"l_shipmode","l_comment", "null"]
csv_format = ds.CsvFileFormat(read_options=csv.ReadOptions(column_names=lineitem_scheme, block_size= 32 * 1024 * 1024), parse_options=csv.ParseOptions(delimiter="|"))
dataset = ds.dataset("s3://TPC",format=csv_format)
s = dataset.to_batches(batch_size=1000000000)
while count < 100:
z = next(s)
```
For our purposes let's just make the TPC dataset consist of hundreds of Parquet files each with one row group. (something that Spark would generate). This script would get somewhere around 1Gbps. If you now do
```
s = dataset.to_batches(batch_size=1000000000, fragment_readahead=16)
```
You can get to 2.5Gbps which is the advertised steady rate cap for this instance type.
Authored-by: Ziheng Wang <zihengw@stanford.edu>
Signed-off-by: Antoine Pitrou <antoine@python.org>
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

@marsupialtail@pitrou@bkietz@westonpace@ursabot