Lakebase에 연결

Structured Streaming을 사용하여 기본 제공되는 배치 처리, 자동 재시도 및 워크스페이스 관리 인증을 통해 Lakebase 또는 외부 PostgreSQL 데이터베이스에 기록할 수 있습니다.

Lakebase 싱크를 사용하는 경우

Lakebase 싱크를 사용해 Lakebase 또는 외부 PostgreSQL 데이터베이스에 저지연 스트리밍 쓰기를 하세요. 이 싱크에서는 일괄 처리, 연결 관리 및 오류 처리를 처리하는 사용자 지정 foreach 함수를 구현할 필요가 없습니다.

일반 사용 사례는 다음과 같습니다.

  • 운영 대시보드 또는 고객 관련 기능을 위해 애플리케이션 데이터베이스를 실시간으로 업데이트합니다.
  • 집계 또는 필터링된 스트리밍 결과와 같이 지속적으로 변경되는 데이터를 트랜잭션 데이터베이스에 동기화합니다.
  • 실시간 모드를 사용하여 1초 미만의 대기 시간이 있는 Lakebase 테이블에 구조적 스트리밍 쿼리의 출력을 씁니다.

레이크베이스에서 레이크하우스의 Delta Lake 테이블로 데이터를 동기화하려면 반대 방향인 Lakebase 변경 데이터 피드를 참조하세요.

요구 사항

  • Databricks 런타임 18 LTS 이상.
    • 외부 PostgreSQL 연결을 사용하려면 Databricks Runtime 19 이상을 사용하고 Custom JDBC on UC Compute 프리뷰에 등록해야 합니다.
    • 인터벌 데이터 타입은 Databricks Runtime 19 이상을 사용해야 합니다.
  • 전용 또는 표준 접근 모드가 있는 클래식 컴퓨트, 또는 노트북이나 작업을 위한 서버리스 컴퓨트가 있습니다. 서버리스 컴퓨트에서는 Trigger.AvailableNow()을(를) 사용하세요. 서버리스 컴퓨팅의 스트리밍을 참조하세요.
  • Lakebase 데이터베이스나 Unity 카탈로그를 외부 PostgreSQL 데이터베이스에 연결하는 경우입니다.

식별자 요구사항

모든 대상에 대해 Databricks는 글자 또는 밑줄로 시작하는 스키마, 테이블, 열, 기본 키 열명 이름을 사용하며, 글자, 숫자, 밑줄만 포함하는 것을 권장합니다. 싱크는 자동으로 Lakebase 테이블을 생성할 때 이러한 요구사항을 강제합니다. 이 요구사항을 충족하지 않는 식별자를 사용하려면 쿼리를 시작하기 전에 대상 테이블을 생성하세요.

데이터베이스에 연결

Lakebase 싱크는 다음 연결 방법을 지원합니다.

Unity 카탈로그에 등록된 Lakebase 테이블

Unity 카탈로그에 등록된 Lakebase 테이블의 경우 커넥터는 자격 증명을 자동으로 관리하고 쿼리를 실행하는 사용자 또는 서비스 주체의 ID를 사용합니다. 테이블이 없으면 커넥터가 테이블을 만듭니다.

Unity 카탈로그에 Lakebase 데이터베이스를 등록하려면 Unity 카탈로그에 Lakebase 데이터베이스 등록을 참조하세요.

Lakebase 테이블에 작성하려면 완전히 제한된 테이블 이름을 가진 메서드를 사용 .toTable() 하세요: catalog.schema.table

Python

(df.writeStream
  .outputMode("update")
  .option("upsertkey", "<primary-key-columns>")  # Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .toTable("<catalog>.<schema>.<table>")
)

Scala

df.writeStream
  .outputMode("update")
  .option("upsertkey", "<primary-key-columns>")  // Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .toTable("<catalog>.<schema>.<table>")

다음 자리 표시자를 바꾸세요:

  • <catalog>.<schema>.<table>: 대상 테이블의 정규화된 이름입니다. Lakebase catalog 데이터베이스를 등록할 때 만든 Unity 카탈로그 카탈로그입니다. Unity 카탈로그에서 Lakebase 데이터베이스 등록을 참조하세요. 테이블이 없으면 커넥터가 테이블을 생성합니다.
  • <primary-key-columns>: 선택 사항입니다. 대상 테이블 기본 키의 모든 열을 쉼표로 구분한 목록입니다. 예: id 또는 user_id,event_type. Upsert 동작 방식을 참고하세요.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: 쿼리가 검사점을 저장하는 Unity 카탈로그 볼륨 경로입니다. 클라우드 개체 스토리지 URI를 사용할 수도 있습니다. 위치는 로컬 디스크가 아니라 쓸 수 있는 스토리지여야 하며 각 스트리밍 쿼리에 고유해야 합니다. 이는 대상 테이블과 독립적입니다. 구조적 스트리밍 검사점을 참조하세요.

batchinterval 및 batchsize와 같은 선택적 구성에 대해서는 PostgreSQL 싱크 옵션을 참조하세요.

Unity 카탈로그에 등록되지 않은 Lakebase 테이블

Unity 카탈로그에 등록되지 않은 Lakebase 테이블의 경우 커넥터는 자격 증명을 자동으로 관리하고 쿼리를 실행하는 사용자 또는 서비스 주체의 ID를 사용합니다. 테이블이 없으면 커넥터가 테이블을 만듭니다.

Lakebase 테이블에 쓰려면 다음 dbtable 옵션과 endpoint 옵션을 사용하세요:

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-columns>")  # Optional
  .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-columns>")  // Optional
  .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>: 선택 사항입니다. 대상 PostgreSQL 데이터베이스의 이름입니다. 기본값은 databricks_postgres입니다. 데이터베이스 관리를 참조하세요.
  • <schema>.<table>: 대상 테이블 형식입니다 schema.table . 스키마를 생략하면 싱크는 public 스키마를 사용합니다. 자동 테이블 생성을 위해서는 글자나 밑줄로 시작하는 식별자를 사용하고, 글자, 숫자, 밑줄만 포함하세요.
  • <primary-key-columns>: 선택 사항입니다. 대상 테이블 기본 키의 모든 열을 쉼표로 구분한 목록입니다. 예: id 또는 user_id,event_type. Upsert 동작 방식을 참고하세요.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: 쿼리가 검사점을 저장하는 Unity 카탈로그 볼륨 경로입니다. 클라우드 개체 스토리지 URI를 사용할 수도 있습니다. 위치는 로컬 디스크가 아니라 쓸 수 있는 스토리지여야 하며 각 스트리밍 쿼리에 고유해야 합니다. 이는 대상 테이블과 독립적입니다. 구조적 스트리밍 검사점을 참조하세요.

batchinterval 및 batchsize와 같은 선택적 구성은 PostgreSQL 싱크 옵션을 참조하세요.

Unity 카탈로그 자격 증명을 사용하는 외부 PostgreSQL

Important

이 기능은 공개 미리보기 단계에 있습니다. 워크스페이스 관리자는 Previews 페이지에서 UC Compute의 커스텀 JDBC 접근 권한을 제어할 수 있습니다. Azure Databricks 미리 보기 관리를 참조하세요.

Unity 카탈로그 연결을 사용해 코드에 자격 증명을 저장하지 않고 외부 PostgreSQL 데이터베이스에 인증하세요. 타겟 테이블은 이미 존재해야 합니다.

연결 생성 유형, POSTGRESQL연결 만들기를 참조하세요. 쿼리를 실행하는 사용자 또는 서비스 주체는 해당 연결에 대해 USE CONNECTION 권한이 있어야 합니다.

PostgreSQL 테이블에 글을 쓰려면 , , database와 dbtable 옵션을 사용databricks.connection하세요:

Python

(df.writeStream
  .format("postgresql")
  .outputMode("update")
  .option("databricks.connection", "<connection-name>")
  .option("database", "<database>")
  .option("dbtable", "<schema>.<table>")
  .option("upsertkey", "<primary-key-columns>")  # Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .start()
)

Scala

df.writeStream
  .format("postgresql")
  .outputMode("update")
  .option("databricks.connection", "<connection-name>")
  .option("database", "<database>")
  .option("dbtable", "<schema>.<table>")
  .option("upsertkey", "<primary-key-columns>")  // Optional
  .option("checkpointLocation", "/Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>")
  .start()

다음 자리 표시자를 바꾸세요:

  • <connection-name>: Unity 카탈로그 연결의 이름입니다.
  • <database>: 대상 PostgreSQL 데이터베이스의 이름입니다.
  • <schema>.<table>: 기존 타겟 테이블 schema.table 형식. 스키마를 생략하면 싱크는 public 스키마를 사용합니다.
  • <primary-key-columns>: 선택 사항입니다. 대상 테이블 기본 키의 모든 열을 쉼표로 구분한 목록입니다. 예: id 또는 user_id,event_type. Upsert 동작 방식을 참고하세요.
  • /Volumes/<catalog>/<schema>/<volume>/<checkpoint-name>: 쿼리가 검사점을 저장하는 Unity 카탈로그 볼륨 경로입니다. 클라우드 개체 스토리지 URI를 사용할 수도 있습니다. 위치는 로컬 디스크가 아니라 쓸 수 있는 스토리지여야 하며 각 스트리밍 쿼리에 고유해야 합니다. 이는 대상 테이블과 독립적입니다. 구조적 스트리밍 검사점을 참조하세요.

PostgreSQL 연결은 항상 TLS를 사용합니다. 인증서 검증은 연결을 생성할 때 선택하는 Unity 카탈로그 연결 설정에 따라 이루어집니다:

  • 신뢰 서버 인증서: 선택되면 연결 은 를 사용하여 sslmode=require서버 인증서를 검증하지 않고 연결을 암호화합니다.
  • 사용자 제공 서버 인증서: 신뢰 서버 인증서가 선택되지 않을 때 사용할 sslmode=verify-full PEM으로 인코딩된 서버 인증서를 제공합니다. 인증서를 제공하지 않으면 연결은 JVM 기본 신뢰 저장소를 사용합니다 sslmode=verify-full .

구성 옵션

싱크에서 인식할 수 없는 옵션 JDBC_STREAMING_SINK_INVALID_OPTIONS에 대한 오류가 발생합니다.

싱크 구성 옵션, 공통 옵션과 각 연결 방법별 옵션은 PostgreSQL 싱크 옵션을 참조하세요.

데이터 유형 매핑

싱크는 각 DataFrame 열이 해당 대상 열과 호환되는지 확인한 후 기존 Lakebase 또는 외부 PostgreSQL 테이블에 기록합니다.

다음 표는 Databricks Runtime 18 LTS 이상에서 지원되는 타입을 포함하고 있습니다:

스파크 유형 자동으로 생성된 Lakebase 테이블 유형 기존 PostgreSQL 테이블에서의 호환 타입
ByteType, ShortType smallint smallint
IntegerType integer integer
LongType bigint bigint
FloatType real real
DoubleType double precision double precision
DecimalType numeric numeric
StringType text varchar, text
VarcharType(n) varchar(n) varchar, text
CharType(n) char(n) char
BinaryType bytea bytea
BooleanType boolean boolean
TimestampType timestamptz timestamptz
TimestampNTZType timestamp timestamp
DateType date date
ArrayType, MapType, StructType, VariantTypeNullType jsonb json, jsonb

다음 표는 Databricks 런타임 19 이상에서 지원하는 타입을 포함하고 있습니다:

스파크 유형 자동으로 생성된 Lakebase 테이블 유형 기존 PostgreSQL 테이블에서의 호환 타입
DayTimeIntervalType, YearMonthIntervalType interval interval

업서트 동작 방식

upsertkey 이 옵션은 대상 테이블의 주요 키 열을 식별합니다. 기존 테이블의 경우, 열 upsertkey 은 테이블의 기본 키와 정확히 일치해야 합니다. 옵션을 생략하면 싱크가 테이블에서 기본 키를 읽습니다. 싱크가 생성하는 Lakebase 테이블의 경우, upsertkey 이 기본 키를 정의합니다. 옵션을 생략하면 싱크가 기본 키 없이 테이블을 생성합니다.

대상 테이블에 기본 키가 있을 때, 싱크는 PostgreSQL의 INSERT INTO ... ON CONFLICT (<primary_key_columns>) DO UPDATE SET ... 문법으로 업서트됩니다. 대상 테이블에 기본 키가 없을 때, 싱크는 삽입을 수행합니다. 쿼리의 출력 모드는 이 동작에 영향을 미치지 않습니다.

모든 기본 키 열은 데이터프레임에 존재해야 하며, 숫자 타입이나 문자열 타입과 같은 비교 가능한 타입을 사용해야 합니다.

성능 튜닝

일괄 처리 및 백프레서

두 조건 중 하나가 충족되면 플러시가 트리거됩니다.

  • 버퍼가 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 No
ProcessingTime Yes No
AvailableNow Yes Yes
Once 예. Deprecated. AvailableNow을 사용합니다. 예. Deprecated. AvailableNow을 사용합니다.

출력 모드

이 표에서는 구조적 스트리밍 출력 모드에 대한 지원을 보여 줍니다.

출력 모드 지원됨
update Yes
append 예. 동작은 update와 동일합니다. 대상 테이블에 기본 키가 있으면 쿼리는 업서트를 수행하고, 그렇지 않으면 삽입을 수행합니다. Upsert 동작 방식을 참고하세요.
complete No

제한 사항

  • Unity 카탈로그 연결을 통해 연결된 외부 PostgreSQL 데이터베이스의 경우, 대상 테이블이 이미 존재해야 합니다. 싱크는 Lakebase에서만 누락된 테이블을 자동으로 생성합니다.
  • 레이크플로우 파이프라인은 지원되지 않습니다.