Important
이 기능은 공개 미리보기 단계에 있습니다.
구조적 스트리밍을 사용하여 기본 제공 일괄 처리, 자동 재시도 및 작업 영역 관리 인증을 사용하여 Lakebase에 씁니다.
Lakebase 싱크를 사용하는 경우
짧은 지연 시간으로 Lakebase에 스트리밍 쓰기를 수행하려면 Lakebase 싱크를 사용하세요. 이 싱크에서는 일괄 처리, 연결 관리 및 오류 처리를 처리하는 사용자 지정 foreachBatch 함수를 구현할 필요가 없습니다.
일반 사용 사례는 다음과 같습니다.
- 운영 대시보드 또는 고객 관련 기능을 위해 애플리케이션 데이터베이스를 실시간으로 업데이트합니다.
- 집계 또는 필터링된 스트리밍 결과와 같이 지속적으로 변경되는 데이터를 트랜잭션 데이터베이스에 동기화합니다.
- 실시간 모드를 사용하여 1초 미만의 대기 시간이 있는 Lakebase 테이블에 구조적 스트리밍 쿼리의 출력을 씁니다.
레이크베이스에서 레이크하우스의 Delta Lake 테이블로 데이터를 동기화하려면 반대 방향인 Lakebase 변경 데이터 피드를 참조하세요.
요구 사항
- Databricks Runtime 18 이상
- 전용 또는 표준 액세스 모드를 사용하는 클래식 컴퓨팅.
- Lakebase 데이터베이스
데이터베이스에 연결
Lakebase 싱크는 다음 연결 방법을 지원합니다.
Unity 카탈로그에 등록된 Lakebase 테이블
Unity 카탈로그에 등록된 Lakebase 테이블의 경우 커넥터는 자격 증명을 자동으로 관리하고 쿼리를 실행하는 사용자 또는 서비스 주체의 ID를 사용합니다. 테이블이 없으면 커넥터가 테이블을 만듭니다.
Unity 카탈로그에 Lakebase 데이터베이스를 등록하려면 Unity 카탈로그에 Lakebase 데이터베이스 등록을 참조하세요.
Lakebase 테이블에 쓰려면 완전 수식된 테이블 이름 .toTable()과 함께 catalog.schema.table 메서드를 사용하세요. 다음 예제에서는 필수 옵션과 선택적 upsertkey 옵션을 보여줍니다.
Python
(df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-column>") # Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
)
Scala
df.writeStream
.outputMode("update")
.option("upsertkey", "<primary-key-column>") // Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.toTable("<catalog>.<schema>.<table>")
다음 자리 표시자를 바꾸세요:
-
<catalog>.<schema>.<table>: 대상 테이블의 정규화된 이름입니다. Lakebasecatalog데이터베이스를 등록할 때 만든 Unity 카탈로그 카탈로그입니다. Unity 카탈로그에서 Lakebase 데이터베이스 등록을 참조하세요. 테이블이 없으면 커넥터가 테이블을 생성합니다. -
<primary-key-column>: 선택 사항입니다. upsert 키를 형성하는 열의 쉼표로 구분된 목록(예iduser_id,event_type: 생략upsertkey하면 싱크는 대상 테이블의 기본 키에서 키를 유추합니다. Upsert 동작을 참조하세요. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: 쿼리가 검사점을 저장하는 Unity 카탈로그 볼륨 경로입니다. 클라우드 개체 스토리지 URI를 사용할 수도 있습니다. 위치는 로컬 디스크가 아니라 쓸 수 있는 스토리지여야 하며 각 스트리밍 쿼리에 고유해야 합니다. 이는 대상 테이블과 독립적입니다. 구조적 스트리밍 검사점을 참조하세요.
같은 선택적 구성은 batchsizebatchinterval구성 옵션을 참조하세요.
Unity 카탈로그에 등록되지 않은 Lakebase 테이블
Unity 카탈로그에 등록되지 않은 Lakebase 테이블의 경우 커넥터는 자격 증명을 자동으로 관리하고 쿼리를 실행하는 사용자 또는 서비스 주체의 ID를 사용합니다. 테이블이 없으면 커넥터가 테이블을 만듭니다.
Lakebase 테이블에 쓰기 위해서는 endpoint 및 dbtable 옵션을 사용합니다. 다음 예제에는 선택 사항 database 및 upsertkey 옵션도 포함되어 있습니다.
Python
(df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") # Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-column>") # Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
)
Scala
df.writeStream
.format("postgresql")
.outputMode("update")
.option("endpoint", "<project-id>.<branch-id>.<endpoint-id>")
.option("database", "<database>") // Optional. Defaults to databricks_postgres.
.option("dbtable", "<schema>.<table>")
.option("upsertkey", "<primary-key-column>") // Optional. Inferred from the table's primary key if omitted.
.option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
.start()
다음 자리 표시자를 바꾸세요:
-
<project-id>.<branch-id>.<endpoint-id>: Lakebase 엔드포인트입니다. [컴퓨팅] 탭의 [ID 가져오기] 메뉴에서 [리소스] 이름에 있는 세 가지 값을 모두 찾습니다. 이 메뉴의 형식projects/<project-id>/branches/<branch-id>/endpoints/<endpoint-id>은 다음과 같습니다. 컴퓨팅 식별자를 참조하세요. -
<database>: 선택 사항입니다. 대상 Postgres 데이터베이스의 이름입니다. 기본값은databricks_postgres입니다. 데이터베이스 관리를 참조하세요. -
<schema>.<table>: 대상 테이블 형식입니다schema.table. 스키마를 생략하면 싱크는public스키마를 사용합니다. 문자 또는 밑줄로 시작하고 문자, 숫자 및 밑줄만 포함하는 간단한 식별자를 사용합니다. 따옴표 붙은 식별자 및 하이픈과 같은 특수 문자는 지원되지 않습니다. -
<primary-key-column>: 선택 사항입니다. upsert 키를 형성하는 열의 쉼표로 구분된 목록(예iduser_id,event_type: 생략upsertkey하면 싱크는 대상 테이블의 기본 키에서 키를 유추합니다. Upsert 동작을 참조하세요. -
/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: 쿼리가 검사점을 저장하는 Unity 카탈로그 볼륨 경로입니다. 클라우드 개체 스토리지 URI를 사용할 수도 있습니다. 위치는 로컬 디스크가 아니라 쓸 수 있는 스토리지여야 하며 각 스트리밍 쿼리에 고유해야 합니다. 이는 대상 테이블과 독립적입니다. 구조적 스트리밍 검사점을 참조하세요.
같은 선택적 구성은 batchsizebatchinterval구성 옵션을 참조하세요.
구성 옵션
싱크에서 인식할 수 없는 옵션 JDBC_STREAMING_SINK_INVALID_OPTIONS에 대한 오류가 발생합니다.
다음 옵션은 모든 연결 메서드에 적용됩니다.
| Key | Default | Description |
|---|---|---|
batchinterval |
100 milliseconds |
Optional. 플러시하기 전에 버퍼에 행을 저장할 최대 시간입니다.
"50 milliseconds"을 예로 들 수 있습니다. |
batchsize |
1000 |
Optional. 각 데이터베이스 트랜잭션에 대한 최대 행 수입니다. |
checkpointLocation |
없음 | Required. Unity 카탈로그 볼륨(/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>)과 같은 검사점 디렉터리의 경로입니다. 각 쿼리에 고유해야 합니다.
구조적 스트리밍 검사점을 참조하세요. |
upsertkey |
없음 | Optional. upsert 키를 형성하는 열 이름의 쉼표로 구분된 목록입니다. 예를 들어 "id" 또는 "user_id,event_type". 지정 upsertkey하는 경우 열이 테이블의 기본 키와 일치해야 합니다. 그렇지 않으면 쿼리가 실패합니다. 이를 생략하면, 싱크는 기본 키를 자동으로 사용합니다. 자세한 내용은 Upsert 동작을 참조하세요. |
Unity 카탈로그에 등록되지 않은 Lakebase 테이블
다음 옵션은 Unity 카탈로그에 등록되지 않은 Lakebase 테이블에 연결할 때 적용됩니다.
| Key | Default | Description |
|---|---|---|
database |
databricks_postgres |
Optional. 대상 PostgreSQL 데이터베이스 이름입니다. |
dbtable |
없음 | Required. 형식의 schema.table 대상 테이블 이름입니다. 스키마를 지정하지 않으면 기본 스키마 값은 .입니다 public. 문자 또는 밑줄로 시작하고 문자, 숫자 및 밑줄만 포함하는 간단한 식별자를 사용합니다. 테이블 또는 스키마 이름을 인용하지 마세요. 하이픈과 같은 특수 문자가 있는 따옴표 붙은 식별자 및 이름은 지원되지 않습니다. |
endpoint |
없음 | Required.
project_id.branch_id 또는 project_id.branch_id.endpoint_id 형식의 Lakebase 엔드포인트.
endpoint_id는 선택 사항입니다. 이를 생략하고 분기에 읽기-쓰기 엔드포인트가 하나만 있는 경우 싱크는 기본값으로 해당 엔드포인트를 선택합니다. |
Upsert 동작
upsert 키가 있고, 그 키가 upsertkey로 지정되었거나 싱크가 테이블의 기본 키에서 추론한 경우, 싱크는 PostgreSQL의 INSERT INTO ... ON CONFLICT (<upsert_key>) DO UPDATE SET ... 구문을 사용해 테이블에 upsert합니다.
upsert 키가 없으면 싱크에서 삽입을 수행합니다. 쿼리의 출력 모드는 upsert 또는 삽입 동작에 영향을 주지 않습니다.
upsertkey 열은 다음과 같아야 합니다:
- DataFrame 열의 비어있지 않은 하위 집합이어야 합니다.
- 대상 테이블의
PRIMARY KEY정확히 일치합니다. 지정한 열이 기본 키와 일치하지 않으면 쿼리가 실패합니다. - 숫자 또는 문자열 형식과 같은 비교 가능한 형식이어야 합니다. 동시 쓰기 중에 데이터베이스 교착 상태를 방지하기 위해 싱크는 각 일괄 처리 내에서 upsert 키로 행을 정렬합니다. Upsert 키는 복합 또는 구조체 형식을 지원하지 않습니다.
열 이름은 PostgreSQL의 기본 방식인 큰따옴표 "로 자동으로 묶이며, 이 방식은 예약어와 대소문자가 혼합된 이름을 처리합니다.
테이블 및 스키마 이름은 문자 또는 밑줄로 시작하고 문자, 숫자 및 밑줄만 포함하는 간단한 식별자를 사용해야 합니다. 싱크는 테이블 또는 스키마 이름에서 따옴표 붙은 식별자 또는 하이픈과 같은 특수 문자를 지원하지 않습니다.
성능 튜닝
일괄 처리 및 백프레서
두 조건 중 하나가 충족되면 플러시가 트리거됩니다.
- 버퍼가
batchsize행에 도달하면 기본값은1000입니다. - 버퍼 유지 기간이 기본값이
batchinterval인100 milliseconds를 초과합니다.
데이터베이스가 들어오는 데이터 속도를 따라갈 수 없는 경우 싱크는 백프레서 업스트림을 원본으로 전파합니다.
대기 시간 및 처리량 지침:
- 실시간 모드의 저지연 워크로드의 경우 플러시되기 전 최대 시간을 더 짧게 보장하려면
batchinterval를 줄이세요. 개념에 대해서는 실시간 모드 개념 을, 코드 예시는 실시간 모드 예 시를 참조하세요. - 처리량이 높은 워크로드의 경우 각 트랜잭션에 대한 오버헤드를 줄이기 위해 증가
batchsize합니다.
연결 동작
싱크는 실행기에서 연결 풀링을 사용합니다. 기본적으로 각 태스크는 하나의 데이터베이스 연결을 사용합니다.
Databricks는 각 연결에 대해 작업의 기본값 1 을 사용하는 것이 좋습니다. 각 연결에 대한 작업 수를 늘리면 연결 경합이 발생하고 높은 처리량 연결에 대한 대기 시간이 증가할 수 있습니다.
연결에 대한 작업의 비율을 구성하려면 Spark 구성을 spark.databricks.sql.streaming.jdbc.tasksPerConnection 설정합니다. 대상 데이터베이스의 연결 수 제한이 낮은 경우 셔플 파티션 수를 줄이거나 spark.databricks.sql.streaming.jdbc.tasksPerConnection를 늘리십시오.
싱크는 연결 오류, 교착 상태 및 속도 제한을 포함하여 임시 JDBC 오류를 자동으로 다시 시도합니다. 싱크가 모든 재시도를 모두 소모하면 쿼리가 실패합니다.
지원되는 트리거 및 출력 모드
Triggers
이 표에서는 구조적 스트리밍 트리거 유형에 대한 지원을 보여 드립니다.
| 트리거 | 지원됨 |
|---|---|
realTime |
Yes |
ProcessingTime |
Yes |
AvailableNow |
Yes |
Once |
Yes |
출력 모드
이 표에서는 구조적 스트리밍 출력 모드에 대한 지원을 보여 줍니다.
| 출력 모드 | 지원됨 |
|---|---|
update |
Yes |
append |
예. 동작은 update와 동일합니다. 대상 테이블에 기본 키가 있으면 쿼리는 업서트를 수행하고, 그렇지 않으면 삽입을 수행합니다.
Upsert 동작을 참조하세요. |
complete |
No |
제한 사항
- 서버리스 컴퓨팅 및 Lakeflow 파이프라인은 지원되지 않습니다.
- Lakebase만 쓰기 대상으로 지원됩니다. 외부 PostgreSQL 호환 데이터베이스는 지원되지 않습니다.