사용자 정의 함수 작성하기

사용자 정의 함수 작성하기 (Write User Defined Functions)

Go 클라이언트는 Go로 작성한 사용자 정의 함수를 등록할 수 있어요. 입력 행을 단일 출력 값에 매핑하는 스칼라 함수와 행 집합을 만들어내는 테이블 함수가 있어요. 해석되지 않는 테이블 이름을 쿼리로 다시 쓰는 리플레이스먼트 스캔도 등록할 수 있어요.

출처: 문서

본문

개요 (Overview)

Go 클라이언트는 Go로 작성한 사용자 정의 함수를 등록할 수 있어요: 입력 행을 단일 출력 값에 매핑하는 스칼라 함수와 행 집합을 만들어내는 테이블 함수. 또한 해석되지 않는 테이블 이름을 쿼리로 다시 쓰는 리플레이스먼트 스캔도 등록할 수 있어요. 스칼라 함수와 테이블 함수는 db.Conn()으로 꺼내온 연결에 등록하고, 리플레이스먼트 스캔은 Connector에 등록해요. 각 경우 함수는 그 후 SQL에서 내장 함수처럼 호출할 수 있어요. 아래 섹션에서 각각을 빌드하고 등록해요.

스칼라 함수 (Scalar Functions)

스칼라 함수는 duckdb.ScalarFunc 인터페이스를 구현하는 Go 타입으로, 두 메서드가 있어요:

  • Config()는 입력 파라미터 타입, 결과 타입, 그리고 가변(variadic) 함수의 경우 가변 파라미터 타입을 선언하는 duckdb.ScalarFuncConfig를 반환해요. 타입은 duckdb.NewTypeInfo()duckdb.NewListInfo() 같은 복합 생성자로 만듭니다.
  • Executor()는 콜백을 담은 duckdb.ScalarFuncExecutor를 반환해요. 가장 단순한 형태는 RowExecutor로, func(values []driver.Value) (any, error)이며 행마다 한 번 호출돼요.

duckdb.RegisterScalarUDF()로 연결에 인스턴스를 등록해요. 다음은 클라이언트의 scalar_udf 예시에서 가져온 my_length(VARCHAR) -> INTEGER를 정의해요:

type varcharLen struct{}

func (*varcharLen) Config() duckdb.ScalarFuncConfig {
    inputType, err := duckdb.NewTypeInfo(duckdb.TYPE_VARCHAR)
    check(err)
    resultType, err := duckdb.NewTypeInfo(duckdb.TYPE_INTEGER)
    check(err)

    return duckdb.ScalarFuncConfig{
        InputTypeInfos: []duckdb.TypeInfo{inputType},
        ResultTypeInfo: resultType,
    }
}

func (*varcharLen) Executor() duckdb.ScalarFuncExecutor {
    return duckdb.ScalarFuncExecutor{
        RowExecutor: func(values []driver.Value) (any, error) {
            return int32(len(values[0].(string))), nil
        },
    }
}

func main() {
    db, err := sql.Open("duckdb", "")
    check(err)
    defer db.Close()

    conn, err := db.Conn(context.Background())
    check(err)
    defer conn.Close()

    var udf *varcharLen
    check(duckdb.RegisterScalarUDF(conn, "my_length", udf))

    var length int32
    row := db.QueryRow(`SELECT my_length('hello world')`)
    check(row.Scan(&length)) // 11
}

등록되면 함수는 RegisterScalarUDF()에 넘긴 이름으로 SQL에서 호출돼요:

SELECT my_length('hello world');

가변 개수의 후행 인자를 받으려면 config에서 VariadicTypeInfo를 설정해요. 이 필드는 duckdb.TypeInfo를 받으므로, 어떤 타입의 인자든 받으려면 duckdb.NewTypeInfo(duckdb.TYPE_ANY)를 사용해요:

anyType, err := duckdb.NewTypeInfo(duckdb.TYPE_ANY)
check(err)

config := duckdb.ScalarFuncConfig{
    VariadicTypeInfo: anyType,
}

하나의 이름 아래 여러 오버로드를 등록하려면(예: VARCHAR를 받는 것과 LIST를 받는 것), 각각을 독자적인 ScalarFunc로 구현하고 duckdb.RegisterScalarUDFSet()으로 함께 등록해요:

check(duckdb.RegisterScalarUDFSet(conn, "my_length", varcharUDF, listUDF))

행 단위가 아니라 한 번에 전체 데이터 청크를 처리하는 콜백을 원하면, executor에서 RowExecutor 대신 ChunkContextExecutor를 설정해요. scalar_udf 예시가 둘 다 보여줘요.

테이블 함수 (Table Functions)

테이블 함수는 행을 만들어내고 FROM 절에서 쿼리돼요. 이는 duckdb.RowTableFunction으로 기술되는데, 그 Config가 인자 타입을 선언하고 BindArguments 콜백이 호출마다 한 번 실행되어 인자를 바인딩하고 table source를 반환해요. table source는 그 메서드들로 실행을 구동해요:

  • ColumnInfos()는 결과 컬럼을 선언하며, 각각은 이름과 TypeInfo를 가진 duckdb.ColumnInfo예요.
  • Init()은 스캔 전에 한 번 실행돼요.
  • FillRow()는 반복 호출되며 호출마다 duckdb.SetRowValue()로 하나의 duckdb.Row를 채우고, 더 많은 행이 있으면 true, 스캔을 끝내려면 false를 반환해요.
  • Cardinality()는 선택적 행 수 추정치를 반환해요.

duckdb.RegisterTableUDF()로 함수를 등록해요. 다음은 클라이언트의 table_udf 예시에서 가져온, 인자까지의 각 정수마다 한 행을 반환하는 increment(BIGINT)를 정의해요:

type incrementTableUDF struct {
    tableSize  int64
    currentRow int64
}

func bindTableUDF(namedArgs map[string]any, args ...any) (duckdb.RowTableSource, error) {
    return &incrementTableUDF{
        currentRow: 0,
        tableSize:  args[0].(int64),
    }, nil
}

func (udf *incrementTableUDF) ColumnInfos() []duckdb.ColumnInfo {
    t, err := duckdb.NewTypeInfo(duckdb.TYPE_BIGINT)
    check(err)
    return []duckdb.ColumnInfo{{Name: "result", T: t}}
}

func (udf *incrementTableUDF) Init() {}

func (udf *incrementTableUDF) FillRow(row duckdb.Row) (bool, error) {
    if udf.currentRow+1 > udf.tableSize {
        return false, nil // 더 이상 행 없음: 스캔 종료
    }
    udf.currentRow++
    err := duckdb.SetRowValue(row, 0, udf.currentRow)
    return true, err
}

func (udf *incrementTableUDF) Cardinality() *duckdb.CardinalityInfo {
    return &duckdb.CardinalityInfo{Cardinality: uint(udf.tableSize), Exact: true}
}

func main() {
    db, err := sql.Open("duckdb", "")
    check(err)
    defer db.Close()

    conn, err := db.Conn(context.Background())
    check(err)
    defer conn.Close()

    t, err := duckdb.NewTypeInfo(duckdb.TYPE_BIGINT)
    check(err)
    udf := duckdb.RowTableFunction{
        Config:        duckdb.TableFunctionConfig{Arguments: []duckdb.TypeInfo{t}},
        BindArguments: bindTableUDF,
    }
    check(duckdb.RegisterTableUDF(conn, "increment", udf))

    rows, err := db.QueryContext(context.Background(), `SELECT * FROM increment(100)`)
    check(err)
    defer rows.Close()
}
SELECT * FROM increment(100);

스캔이 여러 스레드에서 실행될 수 있는 함수라면 duckdb.ParallelRowTableSource를 구현하고 대신 duckdb.ParallelRowTableFunction을 등록해요. 그 Init()MaxThreads를 설정하는 duckdb.ParallelTableSourceInfo를 반환하고, NewLocalState() 메서드가 스레드별 상태를 제공하며, FillRow()는 그 상태를 받아 각 스레드가 자신의 행 범위를 차지하고 출력할 수 있게 해요. 클라이언트의 table_udf_parallel 예시이 전체 패턴을 보여줘요.

리플레이스먼트 스캔 (Replacement Scans)

리플레이스먼트 스캔은 DuckDB가 해석할 수 없는 테이블 이름을 참조하는 쿼리를 가로채서 테이블 함수 호출로 다시 씁니다. duckdb.RegisterReplacementScan()으로 Connector에 등록하며, duckdb.ReplacementScanCallback 타입의 콜백을 넘겨요. 콜백은 해석되지 않은 테이블 이름을 받고 대체 함수 이름과 그 인자를 반환해요:

connector, err := duckdb.NewConnector("", nil)
check(err)
defer connector.Close()

duckdb.RegisterReplacementScan(connector, func(tableName string) (string, []any, error) {
    // 모든 맨 해석되지 않은 테이블 이름을 그 이름의 Parquet 파일로 취급.
    return "read_parquet", []any{tableName + ".parquet"}, nil
})

db := sql.OpenDB(connector)
// SELECT * FROM foo 는 이제 read_parquet('foo.parquet')에서 읽음.

콜백에서 에러를 반환하면 쿼리 에러로 표면화돼요. 이는 DuckDB 클라이언트가 외부 데이터 소스를 일반 테이블처럼 보이게 만드는 메커니즘이에요.

더 알아보기 (Further Reading)

  • Run Queries — SQL에서 등록된 함수 호출하기.
  • Handle Results — Arrow 스트림을 쿼리 가능한 뷰로 등록하기(외부 데이터를 노출하는 관련 방법).
  • Connect — 함수가 등록되는 연결과 리플레이스먼트 스캔용 커넥터 꺼내기.