# Apache Flink 테이블 간 Inner Join 구현: CRUD 작업 지원
Apache Flink 기반 실시간 데이터 마트에서 스트림 간 inner join은 Debezium에서 발생하는 모든 CRUD 작업을 처리해야 합니다. 이전 접근 방식은 삭제와 업데이트를 무시하여 일대다 관계에서 잘못된 결과를 초래했습니다. 새로운 구현은 작업을 파싱하고 상태를 적절히 관리합니다.
Debezium 메시지 구조가 로직을 정의합니다: op 필드가 작업 유형을 나타내고, before와 after는 변경 전후 데이터를 포함합니다.
Debezium 메시지 파싱
Debezium은 Kafka에 고정 구조의 이벤트를 생성합니다:
{
"op": "(c|r|u|d)",
"source": { ... },
"ts_ms": "...",
"before": [Data, null],
"after": [Data, null]
}
필드 의미:
c(생성):before=null,after=Datar(읽기/스냅샷):before=null,after=Datau(업데이트):before=Data,after=Datad(삭제):before=Data,after=null
매퍼는 op에 따라 현재 데이터를 추출합니다.
데이터 모델 매핑
Flink의 Row를 역직렬화하는 Domain 클래스는 작업 유형을 고려합니다:
public class Domain implements Serializable {
public Integer id;
public Integer user_id;
public String domain_name;
public boolean delete;
public static Domain fromRow(Row row) {
Character op = (Character) row.getField(1);
if (op == null) {
throw new IllegalStateException("Never should happen, if Debezium feels fine");
}
Row domain = (Row) row.getField(op == 'd' ? 0 : 2);
Integer id = (Integer) domain.getField(0);
Integer user_id = (Integer) domain.getField(1);
String domain_name = (String) domain.getField(2); // ispravleno on String
return new Domain(id, user_id, domain_name, op == 'd');
}
}
데이터 선택 로직:
c/r의 경우 —after사용.d의 경우 — 항상before.u의 경우 —after(다음 섹션 상세 설명).
User 모델도 동일합니다.
Inner Join을 위한 KeyedCoProcessFunction
InnerJoinFunction의 핵심 로직은 키(user_id)별로 관련 레코드를 상태에 저장합니다:
- 현재 사용자용
ValueState<User>. - 사용자 도메인용
MapState<Integer, Domain>.
User 도착 시:
- 상태 업데이트.
- 모든 도메인을 반복하며 각 Output 생성.
Domain 도착 시:
- MapState에 저장.
- 사용자 존재 시 Output 생성.
전체 코드:
import org.apache.flink.api.common.functions.OpenContext;
import org.apache.flink.api.common.state.*;
import org.apache.flink.streaming.api.functions.co.KeyedCoProcessFunction;
import org.apache.flink.util.Collector;
import java.io.Serializable;
public class InnerJoinFunction extends KeyedCoProcessFunction<Integer, User, Domain, InnerJoinFunction.Output> {
private MapState<Integer, Domain> domainsState;
private ValueState<User> usersState;
@Override
public void processElement1(final User user, final Context ctx, final Collector<InnerJoinFunction.Output> out) throws Exception {
usersState.update(user);
for (final Domain domain : domainsState.values()) {
out.collect(new InnerJoinFunction.Output(
user.id,
user.firstname,
user.lastname,
domain.domain_name,
user.delete || domain.delete
));
}
}
@Override
public void processElement2(final Domain domain, final Context ctx, final Collector<InnerJoinFunction.Output> out) throws Exception {
domainsState.put(domain.id, domain);
final User user = usersState.value();
if (user != null) {
out.collect(new InnerJoinFunction.Output(
user.id,
user.firstname,
user.lastname,
domain.domain_name,
user.delete || domain.delete
));
}
}
@Override
public void open(OpenContext openContext) throws Exception {
var usersStateDescriptor = new ValueStateDescriptor<>(
"users",
User.class
);
var domainsStateDescriptor = new MapStateDescriptor<>(
"domains",
Integer.class,
Domain.class
);
usersState = getRuntimeContext().getState(usersStateDescriptor);
domainsState = getRuntimeContext().getMapState(domainsStateDescriptor);
super.open(openContext);
}
public static class Output implements Serializable {
public Integer user_id;
public String firstname;
public String lastname;
public String domain_name;
public boolean delete;
// getters and setters omitted
}
}
구현 제한 사항
현재 접근 방식은 CRUD를 올바르게 처리하지만, 업데이트 중 join 키(user_id)가 변경되면 실패합니다. 상태는 키에 로컬로 바인딩되어 새 키 파티션으로 재전송되지 않습니다.
해결책은 다음 부분에 있습니다.
주요 포인트:
- Debezium의
op파싱이before/after선택을 결정합니다. - 모델의
delete플래그가 데이터 마트에서 제거를 보장합니다. - 일대다: MapState 반복으로 모든 조합 생성.
- 대용량 성능을 위해 상태 백엔드가 중요합니다.
- join 키 변경 시 특별 처리 필요.
— Editorial Team
아직 댓글이 없습니다.