Tyler Van Hensbergen
11/16/2024, 8:30 PMcredentials_provider function:
def get_credentials() -> daft.io.S3Credentials:
session = boto3.Session()
creds = session.get_credentials()
return daft.io.S3Credentials(
key_id=creds.access_key,
access_key=creds.secret_key,
session_token=creds.token,
expiry=datetime.datetime.now(datetime.timezone.utc) + datetime.timedelta(hours=1),
)
s3_config = S3Config(
credentials_provider=get_credentials,
)
io_config = IOConfig(s3=s3_config)
source = daft.read_parquet(
"<s3://noetik-datalake-prod/branches/main/warehouse/curated__cosmx_cells/*>",
io_config=io_config,
)
io_config = IOConfig(s3=s3_config)
source = daft.read_parquet(
"<s3://noetik-datalake-prod/branches/main/warehouse/curated__cosmx_cells/*>",
io_config=io_config,
)
... but that gives me an opaque serialization error from rust.
The following works initially:
session = boto3.Session()
creds = session.get_credentials()
s3_config = S3Config(
key_id=creds.access_key,
access_key=creds.secret_key,
session_token=creds.token,
)
io_config = IOConfig(s3=s3_config)
source = daft.read_parquet(
"<s3://noetik-datalake-prod/branches/main/warehouse/curated__cosmx_cells/*>",
io_config=io_config,
)
but then the credentials eventually expire so it isn't appropriate for long running jobs that need write access over a long period of time.jay
11/16/2024, 8:46 PMTyler Van Hensbergen
11/16/2024, 8:47 PMSammy Sidhu
11/16/2024, 9:02 PMTyler Van Hensbergen
11/16/2024, 9:03 PM(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) thread '<unnamed>' panicked at src/daft-scan/src/python.rs:480:5: [repeated 30x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) pyo3_runtime.PanicException: called `Result::unwrap()` on an `Err` value: Custom("AttributeError: type object 'S3Credentials' has no attribute 'expiry'") [repeated 90x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) note: run with `RUST_BACKTRACE=1` environment variable to display a backtrace [repeated 30x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) Traceback (most recent call last): [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) File "python/ray/_raylet.pyx", line 2254, in ray._raylet.task_execution_handler [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) File "python/ray/_raylet.pyx", line 2150, in ray._raylet.execute_task_with_cancellation_handler [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) File "python/ray/_raylet.pyx", line 1839, in ray._raylet.execute_task [repeated 300x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 460, in deserialize_objects [repeated 120x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) return context.deserialize_objects(data_metadata_pairs, object_refs) [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) obj = self._deserialize_object(data, metadata, object_ref) [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 317, in _deserialize_object [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) return self._deserialize_msgpack_data(data, metadata_fields) [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 272, in _deserialize_msgpack_data [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) python_objects = self._deserialize_pickle5_data(pickle5_data) [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 262, in _deserialize_pickle5_data [repeated 60x across cluster]
(ScanWithTask-Project-FanoutHash [Stage:2] pid=3867337) obj = pickle.loads(in_band) [repeated 60x across cluster]jay
11/17/2024, 8:59 PMjay
11/17/2024, 9:17 PMassume-role-with-web-identity (which is the intended behavior for EKS pod IAMs I believe)
import boto3
sess = boto3.session.Session()
sess.get_credentials()
creds = sess.get_credentials()
print(creds.method)
>>> "assume-role-with-web-identity"
Using Daft (local mode), I turned on debug logging and was able to confirm that Daft is correctly detecting the appropriate method to use as the WebIdentityToken. Redacted logs:
import logging
import daft
logging.basicConfig(level=logging.DEBUG)
df = daft.read_csv("<s3://daft-public-data/melbourne-airbnb/melbourne_airbnb.csv>")
...
DEBUG:aws_config.default_provider.credentials:provide_credentials; provider=default_chain
DEBUG:aws_config.meta.credentials.chain:load_credentials; provider=Environment
DEBUG:aws_config.meta.credentials.chain:provider in chain did not provide credentials provider=Environment context=the credential provider was not enabled: environment variable not set (CredentialsNotLoaded(CredentialsNotLoaded { source: "environment variable not set" }))
DEBUG:aws_config.meta.credentials.chain:load_credentials; provider=Profile
DEBUG:aws_config.meta.credentials.chain:provider in chain did not provide credentials provider=Profile context=the credential provider was not enabled: No profiles were defined (CredentialsNotLoaded(CredentialsNotLoaded { source: NoProfilesDefined }))
DEBUG:aws_config.meta.credentials.chain:load_credentials; provider=WebIdentityToken
DEBUG:aws_smithy_client:send_operation;
DEBUG:aws_smithy_client:send_operation; operation="AssumeRoleWithWebIdentity"
DEBUG:aws_smithy_client:send_operation; service="sts"
DEBUG:aws_smithy_http_tower.map_request:map_request; name="resolve_endpoint"
DEBUG:aws_smithy_http_tower.map_request:map_request; name="resolve_endpoint"
DEBUG:aws_smithy_http_tower.map_request:map_request; name="generate_user_agent"
DEBUG:aws_smithy_http_tower.map_request:async_map_request; name="retrieve_credentials"
INFO:tracing.span:lazy_load_credentials;
DEBUG:aws_credential_types.cache.lazy_caching:loaded credentials
INFO:aws_http.auth:credentials cache returned CredentialsNotLoaded, ignoring
...
DEBUG:aws_sdk_sts.operation.assume_role_with_web_identity:request_id=Some("***")
DEBUG:aws_smithy_client:send_operation; status="ok"
DEBUG:aws_config.meta.credentials.chain:loaded credentials provider=WebIdentityToken
INFO:aws_credential_types.cache.lazy_caching:credentials cache miss occurred; added new AWS credentials (took 77.98997ms)
...
Our credentials chain seems to work well with EKS pod IAMs on my setup, but please let us know if you are seeing any unexpected behavior here so we can dig in further!jay
11/17/2024, 9:22 PMTyler Van Hensbergen
11/18/2024, 1:12 AMTyler Van Hensbergen
11/18/2024, 1:15 AMjay
11/18/2024, 1:16 AMjay
11/18/2024, 10:04 PM$AWS_CONTAINER_CREDENTIALS_FULL_URI that was provisioned to the pod actually refers to <http://169.254.170.23/v1/credentials> which is a remote serviceTyler Van Hensbergen
11/18/2024, 10:04 PMTyler Van Hensbergen
11/18/2024, 11:02 PMdaft.S3Credentials constructor. If I remove that key, I get an error about access key not being valid either:
(ScanWithTask-LocalLimit [Stage:2] pid=359147) thread '<unnamed>' panicked at src/daft-scan/src/python.rs:455:5:
(ScanWithTask-LocalLimit [Stage:2] pid=359147) called `Result::unwrap()` on an `Err` value: Custom("AttributeError: type object 'S3Credentials' has no attribute 'access_key'")
(ScanWithTask-LocalLimit [Stage:2] pid=359147) note: run with `RUST_BACKTRACE=1` environment variable to display a backtrace
(ScanWithTask-LocalLimit [Stage:2] pid=359147) Traceback (most recent call last):
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 2254, in ray._raylet.task_execution_handler
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 2150, in ray._raylet.execute_task_with_cancellation_handler
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1805, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1806, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1809, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1837, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1839, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/worker.py", line 841, in deserialize_objects
(ScanWithTask-LocalLimit [Stage:2] pid=359147) return context.deserialize_objects(data_metadata_pairs, object_refs)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 460, in deserialize_objects
(ScanWithTask-LocalLimit [Stage:2] pid=359147) obj = self._deserialize_object(data, metadata, object_ref)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 317, in _deserialize_object
(ScanWithTask-LocalLimit [Stage:2] pid=359147) return self._deserialize_msgpack_data(data, metadata_fields)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 272, in _deserialize_msgpack_data
(ScanWithTask-LocalLimit [Stage:2] pid=359147) python_objects = self._deserialize_pickle5_data(pickle5_data)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 262, in _deserialize_pickle5_data
(ScanWithTask-LocalLimit [Stage:2] pid=359147) obj = pickle.loads(in_band)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) pyo3_runtime.PanicException: called `Result::unwrap()` on an `Err` value: Custom("AttributeError: type object 'S3Credentials' has no attribute 'access_key'")
(ScanWithTask-LocalLimit [Stage:2] pid=359147) Exception ignored in: 'ray._raylet.task_execution_handler'
(ScanWithTask-LocalLimit [Stage:2] pid=359147) Traceback (most recent call last):
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 2254, in ray._raylet.task_execution_handler
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 2150, in ray._raylet.execute_task_with_cancellation_handler
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1805, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1806, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1809, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1837, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "python/ray/_raylet.pyx", line 1839, in ray._raylet.execute_task
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/worker.py", line 841, in deserialize_objects
(ScanWithTask-LocalLimit [Stage:2] pid=359147) return context.deserialize_objects(data_metadata_pairs, object_refs)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 460, in deserialize_objects
(ScanWithTask-LocalLimit [Stage:2] pid=359147) obj = self._deserialize_object(data, metadata, object_ref)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 317, in _deserialize_object
(ScanWithTask-LocalLimit [Stage:2] pid=359147) return self._deserialize_msgpack_data(data, metadata_fields)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 272, in _deserialize_msgpack_data
(ScanWithTask-LocalLimit [Stage:2] pid=359147) python_objects = self._deserialize_pickle5_data(pickle5_data)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) File "/home/noetik/twm/.venv/lib/python3.10/site-packages/ray/_private/serialization.py", line 262, in _deserialize_pickle5_data
(ScanWithTask-LocalLimit [Stage:2] pid=359147) obj = pickle.loads(in_band)
(ScanWithTask-LocalLimit [Stage:2] pid=359147) pyo3_runtime.PanicException: called `Result::unwrap()` on an `Err` value: Custom("AttributeError: type object 'S3Credentials' has no attribute 'access_key'")
(ScanWithTask-LocalLimit [Stage:2] pid=359147) [2024-11-18 23:00:04,425 C 359147 359147] <http://task_receiver.cc:213|task_receiver.cc:213>: Check failed: objects_validjay
11/18/2024, 11:02 PMjay
11/18/2024, 11:03 PMTyler Van Hensbergen
11/18/2024, 11:03 PMTyler Van Hensbergen
11/18/2024, 11:49 PMjay
11/19/2024, 12:53 AMTyler Van Hensbergen
11/19/2024, 12:57 AM0.3.13 for daft and 2.38.0 for rayTyler Van Hensbergen
11/19/2024, 1:21 AMKevin Wang
11/19/2024, 1:41 AMRUST_BACKTRACE=1 and share the output error?Tyler Van Hensbergen
11/19/2024, 4:15 AMKevin Wang
11/19/2024, 4:15 AMTyler Van Hensbergen
11/19/2024, 4:22 AMTyler Van Hensbergen
11/19/2024, 5:37 PMKevin Wang
11/19/2024, 9:37 PMKevin Wang
11/19/2024, 10:04 PMget_credentials (the part -> daft.io.S3Credentials). That fixes it for me. I think for some reason the S3Credentials type is not yet fully created by the time the function is deserialized, and so it fails.Kevin Wang
11/19/2024, 10:04 PMTyler Van Hensbergen
11/19/2024, 10:18 PMjay
11/19/2024, 10:19 PMTyler Van Hensbergen
11/19/2024, 10:20 PMTyler Van Hensbergen
11/19/2024, 11:59 PMRayTaskError(RaySystemError): ray::ReduceMerge-WriteFile [Stage:107]() (pid=1818960, ip=10.0.175.38)
At least one of the input arguments for this task could not be computed:
ray.exceptions.RaySystemError: System error: DaftError::SerdeJsonError invalid type: sequence, expected a byte array containing the pickled partition bytes at line 1 column 175
traceback: Traceback (most recent call last):
File "/home/noetik/twm/.venv/lib/python3.10/site-packages/daft/io/config.py", line 8, in _io_config_from_json
return IOConfig.from_json(io_config_json)
daft.exceptions.DaftCoreException: DaftError::SerdeJsonError invalid type: sequence, expected a byte array containing the pickled partition bytes at line 1 column 175Kevin Wang
11/20/2024, 12:19 AMTyler Van Hensbergen
11/20/2024, 12:23 AM