DataFusion CLI

本地文件/目录

师成师成· 更新于 2026-09-28· 阅读 26 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

本地文件 / 目录

只需将文件名、目录名或远程位置用单引号 ' 包裹起来,即可直接对其执行查询,如以下示例所示。

创建一个用于查询的 CSV 文件。

$ echo "a,b" > data.csv
$ echo "1,2" >> data.csv

查询该单个文件(CLI 还支持 parquet、压缩 CSV、avro、json 等格式)。

$ datafusion-cli
DataFusion CLI v17.0.0
> select * from 'data.csv';
+---+---+
| a | b |
+---+---+
| 1 | 2 |
+---+---+
1 row in set. Query took 0.007 seconds.

你也可以查询具有兼容架构的文件目录:

$ ls data_dir/
data.csv   data2.csv
$ datafusion-cli
DataFusion CLI v16.0.0
> select * from 'data_dir';
+---+---+
| a | b |
+---+---+
| 3 | 4 |
| 1 | 2 |
+---+---+
2 rows in set. Query took 0.007 seconds.

远程文件 / 目录

你也可以直接查询 DataFusion 支持的任何远程位置,无需先将该位置注册为表。例如,要通过 HTTP(S) 读取远程 Parquet 文件,可以使用如下方式:

select count(*) from 'https://datasets.clickhouse.com/hits_compatible/athena_partitioned/hits_1.parquet'
+----------+
| COUNT(*) |
+----------+
| 1000000  |
+----------+
1 row in set. Query took 0.595 seconds.

要从 AWS S3 或 GCS 读取数据,请使用 s3 或 gs 作为协议前缀。例如,要读取名为 my-data-bucket 的 S3 存储桶中的文件,请使用 URL s3://my-data-bucket,并将相关的访问凭据设置为环境变量(例如,对于 AWS S3,可以使用 AWS_ACCESS_KEY_ID 和 AWS_SECRET_ACCESS_KEY)。

> select count(*) from 's3://altinity-clickhouse-data/nyc_taxi_rides/data/tripdata_parquet/';
+------------+
| count(*)   |
+------------+
| 1310903963 |
+------------+

有关更多配置选项,请参阅下文的 CREATE EXTERNAL TABLE 一节。

CREATE EXTERNAL TABLE

也可以通过 CREATE EXTERNAL TABLE 创建由文件或远程位置支撑的表,如下所示。请注意,DataFusion 不支持在文件路径中使用通配符(例如 *);相反,可以直接指定目录路径,以读取该目录中所有兼容的文件。

例如,要创建一个由名为 hits.parquet 的本地 parquet 文件支撑的 hits 表:

CREATE EXTERNAL TABLE hits
STORED AS PARQUET
LOCATION 'hits.parquet';

通过 HTTP(S) 创建由远程 Parquet 文件支持的 hits 表:

CREATE EXTERNAL TABLE hits
STORED AS PARQUET
LOCATION 'https://datasets.clickhouse.com/hits_compatible/athena_partitioned/hits_1.parquet';

在这两种情况下,hits 现在都可以像普通表一样进行查询:

select count(*) from hits;
+----------+
| COUNT(*) |
+----------+
| 1000000  |
+----------+
1 row in set. Query took 0.344 seconds.

从标准输入读取

在类 Unix 系统上,你可以将数据通过管道传入 CLI,并通过将 LOCATION 指向 /dev/stdin 伪文件来进行查询:

$ cat hits.csv | datafusion-cli -c "
CREATE EXTERNAL TABLE hits STORED AS CSV LOCATION '/dev/stdin' OPTIONS ('format.has_header' 'true');
SELECT count(*) FROM hits;"

这适用于 CSV、JSON 和 Parquet。由于标准输入不可寻址(而 Parquet 将其元数据存储在文件末尾),CLI 会先把全部输入缓冲到内存中再进行查询,因此数据必须能够放入内存。标准输入只会被读取一次:缓冲的内容会在同一会话中被后续所有以 /dev/stdin 为来源的表复用。这些表必须声明与第一个表相同的 STORED AS 格式;格式不一致会报错并被拒绝。

SQL 必须通过 -c/--command 或 -f/--file 传入,这样才能把标准输入让出来用于传输数据。在交互式 shell 中(以及 SQL 通过管道传给 CLI 但未使用 -c/-f 时),标准输入承载的是 SQL 本身,此时 LOCATION '/dev/stdin' 会返回错误。

为什么不支持通配符

虽然通配符(例如 *.parquet 或 **/*.parquet)在某些情况下可能对本地文件系统有效,但 DataFusion CLI 并不支持它们。这是因为通配符并非在所有存储后端(例如 S3、GCS)上都普遍适用。DataFusion 更希望用户指定目录路径,它会自动读取该目录下所有兼容的文件。

例如,以下用法不受支持:

CREATE EXTERNAL TABLE test (
    message TEXT,
    day DATE
)
STORED AS PARQUET
LOCATION 'gs://bucket/*.parquet';

请改用:

CREATE EXTERNAL TABLE test (
    message TEXT,
    day DATE
)
STORED AS PARQUET
LOCATION 'gs://bucket/my_table/';

当指定的目录路径具有与 Hive 兼容的分区结构时,DataFusion CLI 默认会自动解析这些 Hive 列及其取值,并将其纳入表的架构与数据中。给定以下远程对象路径:

gs://bucket/my_table/a=1/b=100/file1.parquet
gs://bucket/my_table/a=2/b=200/file2.parquet

my_table 可以按 Hive 列进行查询与过滤:

CREATE EXTERNAL TABLE my_table
STORED AS PARQUET
LOCATION 'gs://bucket/my_table/';

SELECT count(*) FROM my_table WHERE b=200;
+----------+
| count(*) |
+----------+
| 1        |
+----------+

格式

Parquet

Parquet 的 schema 信息会自动推导得出。

注册单个 Parquet 文件数据源

CREATE EXTERNAL TABLE taxi
STORED AS PARQUET
LOCATION '/mnt/nyctaxi/tripdata.parquet';

注册单个文件夹作为 parquet 数据源。注意:文件夹内的所有文件必须是有效的 parquet 文件,且 schema 兼容。

:::note

路径必须以斜杠 / 结尾

路径必须以 / 结尾,否则 DataFusion 会将该路径视为文件而非目录。

:::

CREATE EXTERNAL TABLE taxi
STORED AS PARQUET
LOCATION '/mnt/nyctaxi/';

Parquet 专属选项

你可以使用 OPTIONS 子句为 parquet 文件指定额外的选项。例如,要读写带有加密设置的 parquet 目录,可以使用:

CREATE EXTERNAL TABLE encrypted_parquet_table
(
double_field double,
float_field float
)
STORED AS PARQUET LOCATION 'pq/' OPTIONS (
    -- encryption
    'format.crypto.file_encryption.encrypt_footer' 'true',
    'format.crypto.file_encryption.footer_key_as_hex' '30313233343536373839303132333435',  -- b"0123456789012345"
    'format.crypto.file_encryption.column_key_as_hex::double_field' '31323334353637383930313233343530', -- b"1234567890123450"
    'format.crypto.file_encryption.column_key_as_hex::float_field' '31323334353637383930313233343531', -- b"1234567890123451"
    -- decryption
    'format.crypto.file_decryption.footer_key_as_hex' '30313233343536373839303132333435', -- b"0123456789012345"
    'format.crypto.file_decryption.column_key_as_hex::double_field' '31323334353637383930313233343530', -- b"1234567890123450"
    'format.crypto.file_decryption.column_key_as_hex::float_field' '31323334353637383930313233343531', -- b"1234567890123451"
);

这里以十六进制格式指定键,因为它们是二进制数据。在 SQL 中可以使用以下方式对这些键进行编码:

select encode('0123456789012345', 'hex');
/*
+----------------------------------------------+
| encode(Utf8("0123456789012345"),Utf8("hex")) |
+----------------------------------------------+
| 30313233343536373839303132333435             |
+----------------------------------------------+
*/

有关可用选项的更多详情,请参阅 DataFusion 中的 Rust TableParquetOptions 文档。

CSV

DataFusion 会自动推断 CSV 的 schema,你也可以显式提供它。

注册一个带表头行的单文件 CSV 数据源:

CREATE EXTERNAL TABLE test
STORED AS CSV
LOCATION '/path/to/aggregate_test_100.csv'
OPTIONS ('has_header' 'true');

使用显式定义的架构注册单个 CSV 文件数据源:

Register a single file csv datasource with explicitly defined schema:使用显式定义的架构注册单个 CSV 文件数据源:

CREATE EXTERNAL TABLE test (
    c1  VARCHAR NOT NULL,
    c2  INT NOT NULL,
    c3  SMALLINT NOT NULL,
    c4  SMALLINT NOT NULL,
    c5  INT NOT NULL,
    c6  BIGINT NOT NULL,
    c7  SMALLINT NOT NULL,
    c8  INT NOT NULL,
    c9  BIGINT NOT NULL,
    c10 VARCHAR NOT NULL,
    c11 FLOAT NOT NULL,
    c12 DOUBLE NOT NULL,
    c13 VARCHAR NOT NULL
)
STORED AS CSV
LOCATION '/path/to/aggregate_test_100.csv';

位置

HTTP(s)

通过 HTTP(S) 读取远程 parquet 文件:

CREATE EXTERNAL TABLE hits
STORED AS PARQUET
LOCATION 'https://datasets.clickhouse.com/hits_compatible/athena_partitioned/hits_1.parquet';

S3

DataFusion CLI 支持通过 CREATE EXTERNAL TABLE 语句和标准的 AWS 配置方法(通过 aws-config AWS SDK crate)来配置 AWS S3。

使用显式凭据从 S3 存储桶中的文件创建外部表:

CREATE EXTERNAL TABLE test
STORED AS PARQUET
OPTIONS(
    'aws.access_key_id' '******',
    'aws.secret_access_key' '******',
    'aws.region' 'us-east-2'
)
LOCATION 's3://bucket/path/file.parquet';

要使用环境变量创建外部表:

$ export AWS_DEFAULT_REGION=us-east-2
$ export AWS_SECRET_ACCESS_KEY=******
$ export AWS_ACCESS_KEY_ID=******

$ datafusion-cli
`datafusion-cli v21.0.0
> create CREATE TABLE test STORED AS PARQUET LOCATION 's3://bucket/path/file.parquet';
0 rows in set. Query took 0.374 seconds.
> select * from test;
+----------+----------+
| column_1 | column_2 |
+----------+----------+
| 1        | 2        |
+----------+----------+
1 row in set. Query took 0.171 seconds.

要从公共 S3 存储桶读取数据而不使用签名,请使用 aws.SKIP_SIGNATURE 选项:

CREATE EXTERNAL TABLE nyc_taxi_rides
STORED AS PARQUET LOCATION 's3://altinity-clickhouse-data/nyc_taxi_rides/data/tripdata_parquet/'
OPTIONS(aws.SKIP_SIGNATURE true);

凭据按以下优先级顺序获取:

  1. 在 CREATE EXTERNAL TABLE 语句的 OPTIONS 子句中显式指定。
  2. 由 aws-config crate 确定(包括 AWS_ACCESS_KEY_ID 和 AWS_SECRET_ACCESS_KEY 等标准环境变量,以及其他 AWS 特有功能)。

如果未指定任何凭据,DataFusion CLI 将对 S3 使用未签名请求,从而允许读取公共存储桶。

支持的配置选项如下:

环境变量配置选项说明
AWS_ACCESS_KEY_IDaws.access_key_id
AWS_SECRET_ACCESS_KEYaws.secret_access_key
AWS_DEFAULT_REGIONaws.region
AWS_ENDPOINTaws.endpoint
AWS_SESSION_TOKENaws.token
AWS_CONTAINER_CREDENTIALS_RELATIVE_URI参见 IAM 角色
AWS_ALLOW_HTTP如果为 "true",允许不使用 TLS 的 HTTP 连接
AWS_SKIP_SIGNATUREaws.skip_signature如果为 "true",则不对请求进行签名
aws.nosignskip_signature 的别名

OSS

阿里云 OSS 数据源必须配置连接凭据

CREATE EXTERNAL TABLE test
STORED AS PARQUET
OPTIONS(
    'aws.access_key_id' '******',
    'aws.secret_access_key' '******',
    'aws.oss.endpoint' 'https://bucket.oss-cn-hangzhou.aliyuncs.com'
)
LOCATION 'oss://bucket/path/file.parquet';

支持的 OPTIONS 有

  • access_key_id
  • secret_access_key
  • endpoint

注意,OSS 的 endpoint 格式必须为:https://{bucket}.{oss-region-endpoint}

COS

腾讯云 COS 数据源必须配置连接凭证。

CREATE EXTERNAL TABLE test
STORED AS PARQUET
OPTIONS(
    'aws.access_key_id' '******',
    'aws.secret_access_key' '******',
    'aws.cos.endpoint' 'https://cos.ap-singapore.myqcloud.com'
)
LOCATION 'cos://bucket/path/file.parquet';

支持的 OPTIONS 包括:

  • access_key_id
  • secret_access_key
  • endpoint

请注意,endpoint 格式的 URL 必须为:https://cos.{cos-region-endpoint}

GCS

Google Cloud Storage 数据源必须配置连接凭据。

例如,要基于 GCS 存储桶中的文件创建外部表:

CREATE EXTERNAL TABLE test
STORED AS PARQUET
OPTIONS(
    'gcp.service_account_path' '/tmp/gcs.json',
)
LOCATION 'gs://bucket/path/file.parquet';

也可以使用环境变量来指定访问信息:

$ export GOOGLE_SERVICE_ACCOUNT=/tmp/gcs.json

$ datafusion-cli
DataFusion CLI v21.0.0
> create external table test stored as parquet location 'gs://bucket/path/file.parquet';
0 rows in set. Query took 0.374 seconds.
> select * from test;
+----------+----------+
| column_1 | column_2 |
+----------+----------+
| 1        | 2        |
+----------+----------+
1 row in set. Query took 0.171 seconds.

支持的配置选项如下:

环境变量配置选项描述
GOOGLE_SERVICE_ACCOUNTgcp.service_account_path服务账号文件的位置
GOOGLE_SERVICE_ACCOUNT_PATHgcp.service_account_path(别名)服务账号文件的位置
SERVICE_ACCOUNTgcp.service_account_path(别名)服务账号文件的位置
GOOGLE_SERVICE_ACCOUNT_KEYgcp.service_account_key序列化为 JSON 的服务账号密钥
GOOGLE_APPLICATION_CREDENTIALSgcp.application_credentials_path应用凭据文件的位置
GOOGLE_BUCKET存储桶名称
GOOGLE_BUCKET_NAME(别名)存储桶名称

评论

登录后参与评论

正在加载评论…