Python: Optimize PyArrow reads - #6673

Merged
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow
Jan 31, 2023
Merged

Python: Optimize PyArrow reads#6673
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow

Conversation

@Fokko

Copy link
Copy Markdown
Contributor

PyArrow is still sluggish when it comes into opening files, and we still see many requests being made to S3.

This PR removes the Dataset, and uses the lower read_table API. Since the read_table API requires to pass in filters in the DNF form, we need to do some additional conversion.

This PR reduces the number of calls from 203 to 165. Requests log:

Before: https://gist.github.com/Fokko/96b4d5b65ec85c95d6e875f6ec19bf50
After: https://gist.github.com/Fokko/282f6b803d83a830465d97f64cf10057

Query used:

frompyiceberg.catalogimportload_catalogcatalog=load_catalog('local')
tbl=catalog.load_table('nyc.taxis')
frompyiceberg.expressionsimportGreaterThanOrEqual, LessThanOrEqual, Andsc=tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()

Logs:

Also, the wall clock time is lower:

iceberggit:(fd-optimize-pyarrow) ✗ timepython3/tmp/vo.pypython3/tmp/vo.py2.38suser2.75ssystem31%cpu16.067totalpython3/tmp/vo.py2.55suser2.57ssystem36%cpu14.097totalpython3/tmp/vo.py2.60suser2.57ssystem32%cpu15.954total
iceberggit:(master) timepython3/tmp/vo.pypython3/tmp/vo.py2.54suser2.71ssystem28%cpu18.499totalpython3/tmp/vo.py2.75suser2.56ssystem24%cpu21.547totalpython3/tmp/vo.py2.75suser2.95ssystem17%cpu32.554total

Keep in mind that these requests are across the great ocean.

PyArrow is still sluggish when it comes into opening files, and we
still see many requests being made to S3.
This PR removes the Dataset, and uses the lower read_table API.
Since the read_table API requires to pass in filters in the DNF
form, we need to do some additional conversion.
This PR reduces the number of calls from 203 to 165 on my test
query:
```python
from pyiceberg.catalog import load_catalog
catalog = load_catalog('local')
tbl = catalog.load_table('nyc.taxis')
from pyiceberg.expressions import GreaterThanOrEqual, LessThanOrEqual, And
sc = tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()
```
Also, clock time is lower:
```python
➜ iceberg git:(fd-optimize-pyarrow) ✗ time python3 /tmp/vo.py
python3 /tmp/vo.py 2.38s user 2.75s system 31% cpu 16.067 total
python3 /tmp/vo.py 2.55s user 2.57s system 36% cpu 14.097 total
python3 /tmp/vo.py 2.60s user 2.57s system 32% cpu 15.954 total
```
```python
➜ iceberg git:(master) time python3 /tmp/vo.py
python3 /tmp/vo.py 2.54s user 2.71s system 28% cpu 18.499 total
python3 /tmp/vo.py 2.75s user 2.56s system 24% cpu 21.547 total
python3 /tmp/vo.py 2.75s user 2.95s system 17% cpu 32.554 total
```
Keep in mind that these request are across the great ocean
@FokkoFokko changed the title Python: Optimize PyArrow reads 🚀🚀🚀Python: Optimize PyArrow readsJan 26, 2023
raise ValueError(f"Missing Iceberg schema in Metadata for file: {path}")

arrow_table = pq.read_table(
source=fout,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🚀

@rdblue

Copy link
Copy Markdown
Contributor

Looks good to me when tests are passing!

@Fokko

Copy link
Copy Markdown
ContributorAuthor

@rdblue thanks for the review. This one is blocked by #6566

@FokkoFokko added this to the Python 0.4.0 release milestone Jan 30, 2023

def expression_to_plain_format(expressions: Tuple[BooleanExpression, ...]) -> List[List[Tuple[str, str, Any]]]:
def expression_to_plain_format(
expressions: Tuple[BooleanExpression, ...], cast_int_to_datetime: bool = False

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

When would this not be set?

@rdblue
rdblue merged commit 9c230f1 into apache:masterJan 31, 2023
@rdblue

Copy link
Copy Markdown
Contributor

Thanks, @Fokko! Nice work.

krvikash pushed a commit to krvikash/iceberg that referenced this pull request Mar 16, 2023
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@Fokko@rdblue
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all \u003cpre\u003e\u003ccode\u003e blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n 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;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks"); } } catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); } })(); (function(){ try { var __m = "github.com"; var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Python: Optimize PyArrow reads - #6673

Merged
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow
Jan 31, 2023
Merged

Python: Optimize PyArrow reads#6673
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow

Conversation

@Fokko

Copy link
Copy Markdown
Contributor

PyArrow is still sluggish when it comes into opening files, and we still see many requests being made to S3.

This PR removes the Dataset, and uses the lower read_table API. Since the read_table API requires to pass in filters in the DNF form, we need to do some additional conversion.

This PR reduces the number of calls from 203 to 165. Requests log:

Before: https://gist.github.com/Fokko/96b4d5b65ec85c95d6e875f6ec19bf50
After: https://gist.github.com/Fokko/282f6b803d83a830465d97f64cf10057

Query used:

frompyiceberg.catalogimportload_catalogcatalog=load_catalog('local')
tbl=catalog.load_table('nyc.taxis')
frompyiceberg.expressionsimportGreaterThanOrEqual, LessThanOrEqual, Andsc=tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()

Logs:

Also, the wall clock time is lower:

iceberggit:(fd-optimize-pyarrow) ✗ timepython3/tmp/vo.pypython3/tmp/vo.py2.38suser2.75ssystem31%cpu16.067totalpython3/tmp/vo.py2.55suser2.57ssystem36%cpu14.097totalpython3/tmp/vo.py2.60suser2.57ssystem32%cpu15.954total
iceberggit:(master) timepython3/tmp/vo.pypython3/tmp/vo.py2.54suser2.71ssystem28%cpu18.499totalpython3/tmp/vo.py2.75suser2.56ssystem24%cpu21.547totalpython3/tmp/vo.py2.75suser2.95ssystem17%cpu32.554total

Keep in mind that these requests are across the great ocean.

PyArrow is still sluggish when it comes into opening files, and we
still see many requests being made to S3.
This PR removes the Dataset, and uses the lower read_table API.
Since the read_table API requires to pass in filters in the DNF
form, we need to do some additional conversion.
This PR reduces the number of calls from 203 to 165 on my test
query:
```python
from pyiceberg.catalog import load_catalog
catalog = load_catalog('local')
tbl = catalog.load_table('nyc.taxis')
from pyiceberg.expressions import GreaterThanOrEqual, LessThanOrEqual, And
sc = tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()
```
Also, clock time is lower:
```python
➜ iceberg git:(fd-optimize-pyarrow) ✗ time python3 /tmp/vo.py
python3 /tmp/vo.py 2.38s user 2.75s system 31% cpu 16.067 total
python3 /tmp/vo.py 2.55s user 2.57s system 36% cpu 14.097 total
python3 /tmp/vo.py 2.60s user 2.57s system 32% cpu 15.954 total
```
```python
➜ iceberg git:(master) time python3 /tmp/vo.py
python3 /tmp/vo.py 2.54s user 2.71s system 28% cpu 18.499 total
python3 /tmp/vo.py 2.75s user 2.56s system 24% cpu 21.547 total
python3 /tmp/vo.py 2.75s user 2.95s system 17% cpu 32.554 total
```
Keep in mind that these request are across the great ocean
@FokkoFokko changed the title Python: Optimize PyArrow reads 🚀🚀🚀Python: Optimize PyArrow readsJan 26, 2023
raise ValueError(f"Missing Iceberg schema in Metadata for file: {path}")

arrow_table = pq.read_table(
source=fout,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🚀

@rdblue

Copy link
Copy Markdown
Contributor

Looks good to me when tests are passing!

@Fokko

Copy link
Copy Markdown
ContributorAuthor

@rdblue thanks for the review. This one is blocked by #6566

@FokkoFokko added this to the Python 0.4.0 release milestone Jan 30, 2023

def expression_to_plain_format(expressions: Tuple[BooleanExpression, ...]) -> List[List[Tuple[str, str, Any]]]:
def expression_to_plain_format(
expressions: Tuple[BooleanExpression, ...], cast_int_to_datetime: bool = False

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

When would this not be set?

@rdblue
rdblue merged commit 9c230f1 into apache:masterJan 31, 2023
@rdblue

Copy link
Copy Markdown
Contributor

Thanks, @Fokko! Nice work.

krvikash pushed a commit to krvikash/iceberg that referenced this pull request Mar 16, 2023
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@Fokko@rdblue
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Python: Optimize PyArrow reads - #6673

Merged
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow
Jan 31, 2023
Merged

Python: Optimize PyArrow reads#6673
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow

Conversation

@Fokko

Copy link
Copy Markdown
Contributor

PyArrow is still sluggish when it comes into opening files, and we still see many requests being made to S3.

This PR removes the Dataset, and uses the lower read_table API. Since the read_table API requires to pass in filters in the DNF form, we need to do some additional conversion.

This PR reduces the number of calls from 203 to 165. Requests log:

Before: https://gist.github.com/Fokko/96b4d5b65ec85c95d6e875f6ec19bf50
After: https://gist.github.com/Fokko/282f6b803d83a830465d97f64cf10057

Query used:

frompyiceberg.catalogimportload_catalogcatalog=load_catalog('local')
tbl=catalog.load_table('nyc.taxis')
frompyiceberg.expressionsimportGreaterThanOrEqual, LessThanOrEqual, Andsc=tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()

Logs:

Also, the wall clock time is lower:

iceberggit:(fd-optimize-pyarrow) ✗ timepython3/tmp/vo.pypython3/tmp/vo.py2.38suser2.75ssystem31%cpu16.067totalpython3/tmp/vo.py2.55suser2.57ssystem36%cpu14.097totalpython3/tmp/vo.py2.60suser2.57ssystem32%cpu15.954total
iceberggit:(master) timepython3/tmp/vo.pypython3/tmp/vo.py2.54suser2.71ssystem28%cpu18.499totalpython3/tmp/vo.py2.75suser2.56ssystem24%cpu21.547totalpython3/tmp/vo.py2.75suser2.95ssystem17%cpu32.554total

Keep in mind that these requests are across the great ocean.

PyArrow is still sluggish when it comes into opening files, and we
still see many requests being made to S3.
This PR removes the Dataset, and uses the lower read_table API.
Since the read_table API requires to pass in filters in the DNF
form, we need to do some additional conversion.
This PR reduces the number of calls from 203 to 165 on my test
query:
```python
from pyiceberg.catalog import load_catalog
catalog = load_catalog('local')
tbl = catalog.load_table('nyc.taxis')
from pyiceberg.expressions import GreaterThanOrEqual, LessThanOrEqual, And
sc = tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()
```
Also, clock time is lower:
```python
➜ iceberg git:(fd-optimize-pyarrow) ✗ time python3 /tmp/vo.py
python3 /tmp/vo.py 2.38s user 2.75s system 31% cpu 16.067 total
python3 /tmp/vo.py 2.55s user 2.57s system 36% cpu 14.097 total
python3 /tmp/vo.py 2.60s user 2.57s system 32% cpu 15.954 total
```
```python
➜ iceberg git:(master) time python3 /tmp/vo.py
python3 /tmp/vo.py 2.54s user 2.71s system 28% cpu 18.499 total
python3 /tmp/vo.py 2.75s user 2.56s system 24% cpu 21.547 total
python3 /tmp/vo.py 2.75s user 2.95s system 17% cpu 32.554 total
```
Keep in mind that these request are across the great ocean
@FokkoFokko changed the title Python: Optimize PyArrow reads 🚀🚀🚀Python: Optimize PyArrow readsJan 26, 2023
raise ValueError(f"Missing Iceberg schema in Metadata for file: {path}")

arrow_table = pq.read_table(
source=fout,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🚀

@rdblue

Copy link
Copy Markdown
Contributor

Looks good to me when tests are passing!

@Fokko

Copy link
Copy Markdown
ContributorAuthor

@rdblue thanks for the review. This one is blocked by #6566

@FokkoFokko added this to the Python 0.4.0 release milestone Jan 30, 2023

def expression_to_plain_format(expressions: Tuple[BooleanExpression, ...]) -> List[List[Tuple[str, str, Any]]]:
def expression_to_plain_format(
expressions: Tuple[BooleanExpression, ...], cast_int_to_datetime: bool = False

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

When would this not be set?

@rdblue
rdblue merged commit 9c230f1 into apache:masterJan 31, 2023
@rdblue

Copy link
Copy Markdown
Contributor

Thanks, @Fokko! Nice work.

krvikash pushed a commit to krvikash/iceberg that referenced this pull request Mar 16, 2023
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

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

Python: Optimize PyArrow reads - #6673

Merged
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow
Jan 31, 2023
Merged

Python: Optimize PyArrow reads#6673
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow

Conversation

@Fokko

Copy link
Copy Markdown
Contributor

PyArrow is still sluggish when it comes into opening files, and we still see many requests being made to S3.

This PR removes the Dataset, and uses the lower read_table API. Since the read_table API requires to pass in filters in the DNF form, we need to do some additional conversion.

This PR reduces the number of calls from 203 to 165. Requests log:

Before: https://gist.github.com/Fokko/96b4d5b65ec85c95d6e875f6ec19bf50
After: https://gist.github.com/Fokko/282f6b803d83a830465d97f64cf10057

Query used:

frompyiceberg.catalogimportload_catalogcatalog=load_catalog('local')
tbl=catalog.load_table('nyc.taxis')
frompyiceberg.expressionsimportGreaterThanOrEqual, LessThanOrEqual, Andsc=tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()

Logs:

Also, the wall clock time is lower:

iceberggit:(fd-optimize-pyarrow) ✗ timepython3/tmp/vo.pypython3/tmp/vo.py2.38suser2.75ssystem31%cpu16.067totalpython3/tmp/vo.py2.55suser2.57ssystem36%cpu14.097totalpython3/tmp/vo.py2.60suser2.57ssystem32%cpu15.954total
iceberggit:(master) timepython3/tmp/vo.pypython3/tmp/vo.py2.54suser2.71ssystem28%cpu18.499totalpython3/tmp/vo.py2.75suser2.56ssystem24%cpu21.547totalpython3/tmp/vo.py2.75suser2.95ssystem17%cpu32.554total

Keep in mind that these requests are across the great ocean.

PyArrow is still sluggish when it comes into opening files, and we
still see many requests being made to S3.
This PR removes the Dataset, and uses the lower read_table API.
Since the read_table API requires to pass in filters in the DNF
form, we need to do some additional conversion.
This PR reduces the number of calls from 203 to 165 on my test
query:
```python
from pyiceberg.catalog import load_catalog
catalog = load_catalog('local')
tbl = catalog.load_table('nyc.taxis')
from pyiceberg.expressions import GreaterThanOrEqual, LessThanOrEqual, And
sc = tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()
```
Also, clock time is lower:
```python
➜ iceberg git:(fd-optimize-pyarrow) ✗ time python3 /tmp/vo.py
python3 /tmp/vo.py 2.38s user 2.75s system 31% cpu 16.067 total
python3 /tmp/vo.py 2.55s user 2.57s system 36% cpu 14.097 total
python3 /tmp/vo.py 2.60s user 2.57s system 32% cpu 15.954 total
```
```python
➜ iceberg git:(master) time python3 /tmp/vo.py
python3 /tmp/vo.py 2.54s user 2.71s system 28% cpu 18.499 total
python3 /tmp/vo.py 2.75s user 2.56s system 24% cpu 21.547 total
python3 /tmp/vo.py 2.75s user 2.95s system 17% cpu 32.554 total
```
Keep in mind that these request are across the great ocean
@FokkoFokko changed the title Python: Optimize PyArrow reads 🚀🚀🚀Python: Optimize PyArrow readsJan 26, 2023
raise ValueError(f"Missing Iceberg schema in Metadata for file: {path}")

arrow_table = pq.read_table(
source=fout,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🚀

@rdblue

Copy link
Copy Markdown
Contributor

Looks good to me when tests are passing!

@Fokko

Copy link
Copy Markdown
ContributorAuthor

@rdblue thanks for the review. This one is blocked by #6566

@FokkoFokko added this to the Python 0.4.0 release milestone Jan 30, 2023

def expression_to_plain_format(expressions: Tuple[BooleanExpression, ...]) -> List[List[Tuple[str, str, Any]]]:
def expression_to_plain_format(
expressions: Tuple[BooleanExpression, ...], cast_int_to_datetime: bool = False

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

When would this not be set?

@rdblue
rdblue merged commit 9c230f1 into apache:masterJan 31, 2023
@rdblue

Copy link
Copy Markdown
Contributor

Thanks, @Fokko! Nice work.

krvikash pushed a commit to krvikash/iceberg that referenced this pull request Mar 16, 2023
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

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

Python: Optimize PyArrow reads - #6673

Merged
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow
Jan 31, 2023
Merged

Python: Optimize PyArrow reads#6673
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow

Conversation

@Fokko

Copy link
Copy Markdown
Contributor

PyArrow is still sluggish when it comes into opening files, and we still see many requests being made to S3.

This PR removes the Dataset, and uses the lower read_table API. Since the read_table API requires to pass in filters in the DNF form, we need to do some additional conversion.

This PR reduces the number of calls from 203 to 165. Requests log:

Before: https://gist.github.com/Fokko/96b4d5b65ec85c95d6e875f6ec19bf50
After: https://gist.github.com/Fokko/282f6b803d83a830465d97f64cf10057

Query used:

frompyiceberg.catalogimportload_catalogcatalog=load_catalog('local')
tbl=catalog.load_table('nyc.taxis')
frompyiceberg.expressionsimportGreaterThanOrEqual, LessThanOrEqual, Andsc=tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()

Logs:

Also, the wall clock time is lower:

iceberggit:(fd-optimize-pyarrow) ✗ timepython3/tmp/vo.pypython3/tmp/vo.py2.38suser2.75ssystem31%cpu16.067totalpython3/tmp/vo.py2.55suser2.57ssystem36%cpu14.097totalpython3/tmp/vo.py2.60suser2.57ssystem32%cpu15.954total
iceberggit:(master) timepython3/tmp/vo.pypython3/tmp/vo.py2.54suser2.71ssystem28%cpu18.499totalpython3/tmp/vo.py2.75suser2.56ssystem24%cpu21.547totalpython3/tmp/vo.py2.75suser2.95ssystem17%cpu32.554total

Keep in mind that these requests are across the great ocean.

PyArrow is still sluggish when it comes into opening files, and we
still see many requests being made to S3.
This PR removes the Dataset, and uses the lower read_table API.
Since the read_table API requires to pass in filters in the DNF
form, we need to do some additional conversion.
This PR reduces the number of calls from 203 to 165 on my test
query:
```python
from pyiceberg.catalog import load_catalog
catalog = load_catalog('local')
tbl = catalog.load_table('nyc.taxis')
from pyiceberg.expressions import GreaterThanOrEqual, LessThanOrEqual, And
sc = tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()
```
Also, clock time is lower:
```python
➜ iceberg git:(fd-optimize-pyarrow) ✗ time python3 /tmp/vo.py
python3 /tmp/vo.py 2.38s user 2.75s system 31% cpu 16.067 total
python3 /tmp/vo.py 2.55s user 2.57s system 36% cpu 14.097 total
python3 /tmp/vo.py 2.60s user 2.57s system 32% cpu 15.954 total
```
```python
➜ iceberg git:(master) time python3 /tmp/vo.py
python3 /tmp/vo.py 2.54s user 2.71s system 28% cpu 18.499 total
python3 /tmp/vo.py 2.75s user 2.56s system 24% cpu 21.547 total
python3 /tmp/vo.py 2.75s user 2.95s system 17% cpu 32.554 total
```
Keep in mind that these request are across the great ocean
@FokkoFokko changed the title Python: Optimize PyArrow reads 🚀🚀🚀Python: Optimize PyArrow readsJan 26, 2023
raise ValueError(f"Missing Iceberg schema in Metadata for file: {path}")

arrow_table = pq.read_table(
source=fout,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🚀

@rdblue

Copy link
Copy Markdown
Contributor

Looks good to me when tests are passing!

@Fokko

Copy link
Copy Markdown
ContributorAuthor

@rdblue thanks for the review. This one is blocked by #6566

@FokkoFokko added this to the Python 0.4.0 release milestone Jan 30, 2023

def expression_to_plain_format(expressions: Tuple[BooleanExpression, ...]) -> List[List[Tuple[str, str, Any]]]:
def expression_to_plain_format(
expressions: Tuple[BooleanExpression, ...], cast_int_to_datetime: bool = False

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

When would this not be set?

@rdblue
rdblue merged commit 9c230f1 into apache:masterJan 31, 2023
@rdblue

Copy link
Copy Markdown
Contributor

Thanks, @Fokko! Nice work.

krvikash pushed a commit to krvikash/iceberg that referenced this pull request Mar 16, 2023
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@Fokko@rdblue
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Python: Optimize PyArrow reads - #6673

Merged
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow
Jan 31, 2023
Merged

Python: Optimize PyArrow reads#6673
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow

Conversation

@Fokko

Copy link
Copy Markdown
Contributor

PyArrow is still sluggish when it comes into opening files, and we still see many requests being made to S3.

This PR removes the Dataset, and uses the lower read_table API. Since the read_table API requires to pass in filters in the DNF form, we need to do some additional conversion.

This PR reduces the number of calls from 203 to 165. Requests log:

Before: https://gist.github.com/Fokko/96b4d5b65ec85c95d6e875f6ec19bf50
After: https://gist.github.com/Fokko/282f6b803d83a830465d97f64cf10057

Query used:

frompyiceberg.catalogimportload_catalogcatalog=load_catalog('local')
tbl=catalog.load_table('nyc.taxis')
frompyiceberg.expressionsimportGreaterThanOrEqual, LessThanOrEqual, Andsc=tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()

Logs:

Also, the wall clock time is lower:

iceberggit:(fd-optimize-pyarrow) ✗ timepython3/tmp/vo.pypython3/tmp/vo.py2.38suser2.75ssystem31%cpu16.067totalpython3/tmp/vo.py2.55suser2.57ssystem36%cpu14.097totalpython3/tmp/vo.py2.60suser2.57ssystem32%cpu15.954total
iceberggit:(master) timepython3/tmp/vo.pypython3/tmp/vo.py2.54suser2.71ssystem28%cpu18.499totalpython3/tmp/vo.py2.75suser2.56ssystem24%cpu21.547totalpython3/tmp/vo.py2.75suser2.95ssystem17%cpu32.554total

Keep in mind that these requests are across the great ocean.

PyArrow is still sluggish when it comes into opening files, and we
still see many requests being made to S3.
This PR removes the Dataset, and uses the lower read_table API.
Since the read_table API requires to pass in filters in the DNF
form, we need to do some additional conversion.
This PR reduces the number of calls from 203 to 165 on my test
query:
```python
from pyiceberg.catalog import load_catalog
catalog = load_catalog('local')
tbl = catalog.load_table('nyc.taxis')
from pyiceberg.expressions import GreaterThanOrEqual, LessThanOrEqual, And
sc = tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()
```
Also, clock time is lower:
```python
➜ iceberg git:(fd-optimize-pyarrow) ✗ time python3 /tmp/vo.py
python3 /tmp/vo.py 2.38s user 2.75s system 31% cpu 16.067 total
python3 /tmp/vo.py 2.55s user 2.57s system 36% cpu 14.097 total
python3 /tmp/vo.py 2.60s user 2.57s system 32% cpu 15.954 total
```
```python
➜ iceberg git:(master) time python3 /tmp/vo.py
python3 /tmp/vo.py 2.54s user 2.71s system 28% cpu 18.499 total
python3 /tmp/vo.py 2.75s user 2.56s system 24% cpu 21.547 total
python3 /tmp/vo.py 2.75s user 2.95s system 17% cpu 32.554 total
```
Keep in mind that these request are across the great ocean
@FokkoFokko changed the title Python: Optimize PyArrow reads 🚀🚀🚀Python: Optimize PyArrow readsJan 26, 2023
raise ValueError(f"Missing Iceberg schema in Metadata for file: {path}")

arrow_table = pq.read_table(
source=fout,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🚀

@rdblue

Copy link
Copy Markdown
Contributor

Looks good to me when tests are passing!

@Fokko

Copy link
Copy Markdown
ContributorAuthor

@rdblue thanks for the review. This one is blocked by #6566

@FokkoFokko added this to the Python 0.4.0 release milestone Jan 30, 2023

def expression_to_plain_format(expressions: Tuple[BooleanExpression, ...]) -> List[List[Tuple[str, str, Any]]]:
def expression_to_plain_format(
expressions: Tuple[BooleanExpression, ...], cast_int_to_datetime: bool = False

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

When would this not be set?

@rdblue
rdblue merged commit 9c230f1 into apache:masterJan 31, 2023
@rdblue

Copy link
Copy Markdown
Contributor

Thanks, @Fokko! Nice work.

krvikash pushed a commit to krvikash/iceberg that referenced this pull request Mar 16, 2023
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@Fokko@rdblue
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Python: Optimize PyArrow reads - #6673

Merged
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow
Jan 31, 2023
Merged

Python: Optimize PyArrow reads#6673
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow

Conversation

@Fokko

Copy link
Copy Markdown
Contributor

PyArrow is still sluggish when it comes into opening files, and we still see many requests being made to S3.

This PR removes the Dataset, and uses the lower read_table API. Since the read_table API requires to pass in filters in the DNF form, we need to do some additional conversion.

This PR reduces the number of calls from 203 to 165. Requests log:

Before: https://gist.github.com/Fokko/96b4d5b65ec85c95d6e875f6ec19bf50
After: https://gist.github.com/Fokko/282f6b803d83a830465d97f64cf10057

Query used:

frompyiceberg.catalogimportload_catalogcatalog=load_catalog('local')
tbl=catalog.load_table('nyc.taxis')
frompyiceberg.expressionsimportGreaterThanOrEqual, LessThanOrEqual, Andsc=tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()

Logs:

Also, the wall clock time is lower:

iceberggit:(fd-optimize-pyarrow) ✗ timepython3/tmp/vo.pypython3/tmp/vo.py2.38suser2.75ssystem31%cpu16.067totalpython3/tmp/vo.py2.55suser2.57ssystem36%cpu14.097totalpython3/tmp/vo.py2.60suser2.57ssystem32%cpu15.954total
iceberggit:(master) timepython3/tmp/vo.pypython3/tmp/vo.py2.54suser2.71ssystem28%cpu18.499totalpython3/tmp/vo.py2.75suser2.56ssystem24%cpu21.547totalpython3/tmp/vo.py2.75suser2.95ssystem17%cpu32.554total

Keep in mind that these requests are across the great ocean.

PyArrow is still sluggish when it comes into opening files, and we
still see many requests being made to S3.
This PR removes the Dataset, and uses the lower read_table API.
Since the read_table API requires to pass in filters in the DNF
form, we need to do some additional conversion.
This PR reduces the number of calls from 203 to 165 on my test
query:
```python
from pyiceberg.catalog import load_catalog
catalog = load_catalog('local')
tbl = catalog.load_table('nyc.taxis')
from pyiceberg.expressions import GreaterThanOrEqual, LessThanOrEqual, And
sc = tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()
```
Also, clock time is lower:
```python
➜ iceberg git:(fd-optimize-pyarrow) ✗ time python3 /tmp/vo.py
python3 /tmp/vo.py 2.38s user 2.75s system 31% cpu 16.067 total
python3 /tmp/vo.py 2.55s user 2.57s system 36% cpu 14.097 total
python3 /tmp/vo.py 2.60s user 2.57s system 32% cpu 15.954 total
```
```python
➜ iceberg git:(master) time python3 /tmp/vo.py
python3 /tmp/vo.py 2.54s user 2.71s system 28% cpu 18.499 total
python3 /tmp/vo.py 2.75s user 2.56s system 24% cpu 21.547 total
python3 /tmp/vo.py 2.75s user 2.95s system 17% cpu 32.554 total
```
Keep in mind that these request are across the great ocean
@FokkoFokko changed the title Python: Optimize PyArrow reads 🚀🚀🚀Python: Optimize PyArrow readsJan 26, 2023
raise ValueError(f"Missing Iceberg schema in Metadata for file: {path}")

arrow_table = pq.read_table(
source=fout,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🚀

@rdblue

Copy link
Copy Markdown
Contributor

Looks good to me when tests are passing!

@Fokko

Copy link
Copy Markdown
ContributorAuthor

@rdblue thanks for the review. This one is blocked by #6566

@FokkoFokko added this to the Python 0.4.0 release milestone Jan 30, 2023

def expression_to_plain_format(expressions: Tuple[BooleanExpression, ...]) -> List[List[Tuple[str, str, Any]]]:
def expression_to_plain_format(
expressions: Tuple[BooleanExpression, ...], cast_int_to_datetime: bool = False

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

When would this not be set?

@rdblue
rdblue merged commit 9c230f1 into apache:masterJan 31, 2023
@rdblue

Copy link
Copy Markdown
Contributor

Thanks, @Fokko! Nice work.

krvikash pushed a commit to krvikash/iceberg that referenced this pull request Mar 16, 2023
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

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

Python: Optimize PyArrow reads - #6673

Merged
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow
Jan 31, 2023
Merged

Python: Optimize PyArrow reads#6673
rdblue merged 4 commits into
apache:masterfrom
Fokko:fd-optimize-pyarrow

Conversation

@Fokko

Copy link
Copy Markdown
Contributor

PyArrow is still sluggish when it comes into opening files, and we still see many requests being made to S3.

This PR removes the Dataset, and uses the lower read_table API. Since the read_table API requires to pass in filters in the DNF form, we need to do some additional conversion.

This PR reduces the number of calls from 203 to 165. Requests log:

Before: https://gist.github.com/Fokko/96b4d5b65ec85c95d6e875f6ec19bf50
After: https://gist.github.com/Fokko/282f6b803d83a830465d97f64cf10057

Query used:

frompyiceberg.catalogimportload_catalogcatalog=load_catalog('local')
tbl=catalog.load_table('nyc.taxis')
frompyiceberg.expressionsimportGreaterThanOrEqual, LessThanOrEqual, Andsc=tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()

Logs:

Also, the wall clock time is lower:

iceberggit:(fd-optimize-pyarrow) ✗ timepython3/tmp/vo.pypython3/tmp/vo.py2.38suser2.75ssystem31%cpu16.067totalpython3/tmp/vo.py2.55suser2.57ssystem36%cpu14.097totalpython3/tmp/vo.py2.60suser2.57ssystem32%cpu15.954total
iceberggit:(master) timepython3/tmp/vo.pypython3/tmp/vo.py2.54suser2.71ssystem28%cpu18.499totalpython3/tmp/vo.py2.75suser2.56ssystem24%cpu21.547totalpython3/tmp/vo.py2.75suser2.95ssystem17%cpu32.554total

Keep in mind that these requests are across the great ocean.

PyArrow is still sluggish when it comes into opening files, and we
still see many requests being made to S3.
This PR removes the Dataset, and uses the lower read_table API.
Since the read_table API requires to pass in filters in the DNF
form, we need to do some additional conversion.
This PR reduces the number of calls from 203 to 165 on my test
query:
```python
from pyiceberg.catalog import load_catalog
catalog = load_catalog('local')
tbl = catalog.load_table('nyc.taxis')
from pyiceberg.expressions import GreaterThanOrEqual, LessThanOrEqual, And
sc = tbl.scan(row_filter=And(
GreaterThanOrEqual("tpep_pickup_datetime", "2022-04-01T00:00:00.000000+00:00"),
LessThanOrEqual("tpep_pickup_datetime", "2022-04-28T00:00:00.000000+00:00"),
)).to_arrow()
```
Also, clock time is lower:
```python
➜ iceberg git:(fd-optimize-pyarrow) ✗ time python3 /tmp/vo.py
python3 /tmp/vo.py 2.38s user 2.75s system 31% cpu 16.067 total
python3 /tmp/vo.py 2.55s user 2.57s system 36% cpu 14.097 total
python3 /tmp/vo.py 2.60s user 2.57s system 32% cpu 15.954 total
```
```python
➜ iceberg git:(master) time python3 /tmp/vo.py
python3 /tmp/vo.py 2.54s user 2.71s system 28% cpu 18.499 total
python3 /tmp/vo.py 2.75s user 2.56s system 24% cpu 21.547 total
python3 /tmp/vo.py 2.75s user 2.95s system 17% cpu 32.554 total
```
Keep in mind that these request are across the great ocean
@FokkoFokko changed the title Python: Optimize PyArrow reads 🚀🚀🚀Python: Optimize PyArrow readsJan 26, 2023
raise ValueError(f"Missing Iceberg schema in Metadata for file: {path}")

arrow_table = pq.read_table(
source=fout,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

🚀

@rdblue

Copy link
Copy Markdown
Contributor

Looks good to me when tests are passing!

@Fokko

Copy link
Copy Markdown
ContributorAuthor

@rdblue thanks for the review. This one is blocked by #6566

@FokkoFokko added this to the Python 0.4.0 release milestone Jan 30, 2023

def expression_to_plain_format(expressions: Tuple[BooleanExpression, ...]) -> List[List[Tuple[str, str, Any]]]:
def expression_to_plain_format(
expressions: Tuple[BooleanExpression, ...], cast_int_to_datetime: bool = False

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

When would this not be set?

@rdblue
rdblue merged commit 9c230f1 into apache:masterJan 31, 2023
@rdblue

Copy link
Copy Markdown
Contributor

Thanks, @Fokko! Nice work.

krvikash pushed a commit to krvikash/iceberg that referenced this pull request Mar 16, 2023
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants

@Fokko@rdblue